use std::sync::Arc;
use std::time::Duration;
use async_trait::async_trait;
use base64::prelude::*;
use serde_json::{Value, json};
use tracing::{error, info, warn};
use crate::error::Problem;
use crate::jobs::{JobHandler, JobOutcome, JobQueue, JobSpec};
use crate::notify::{CertificateIssuedData, NotifyEvent};
use crate::sqlite::job::Job;
use crate::sqlite::order::Order;
use crate::sqlite::upstream_order::UpstreamOrder;
use super::client::{Signer, UpstreamError};
use super::wire::{UpstreamAuthzView, UpstreamOrderView};
use super::{ChallengeStrategy, Inner, dns01, http01};
pub const RELAY_JOB_KIND: &str = "signer_relay_issue";
#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
#[error("{}", self.reason())]
pub(super) enum RelayFailure {
Retryable(String),
Permanent(String),
}
impl RelayFailure {
fn reason(&self) -> &str {
match self {
RelayFailure::Retryable(reason) | RelayFailure::Permanent(reason) => reason,
}
}
}
pub(super) fn classify(error: &UpstreamError) -> RelayFailure {
let reason = error.to_string();
match error {
UpstreamError::Transport(_) | UpstreamError::Url(_) => RelayFailure::Retryable(reason),
UpstreamError::Protocol(_) => RelayFailure::Retryable(reason),
UpstreamError::Problem { status, .. } if *status >= 500 || *status == 429 => {
RelayFailure::Retryable(reason)
}
UpstreamError::Problem { .. } => RelayFailure::Permanent(reason),
UpstreamError::Jws(_) => RelayFailure::Permanent(reason),
}
}
pub struct RelayJob(pub(super) Arc<Inner>);
#[async_trait]
impl JobHandler for RelayJob {
fn kind(&self) -> &'static str {
RELAY_JOB_KIND
}
fn lease(&self) -> Option<Duration> {
Some(self.0.poll.timeout)
}
async fn run(&self, job: &Job) -> JobOutcome {
let inner = &self.0;
let Some(order_id) = job.payload.get("order_id").and_then(Value::as_str) else {
return JobOutcome::Failed("the job payload names no order".to_string());
};
let mapping = match UpstreamOrder::find_by_order_id(order_id, &inner.database).await {
Ok(Some(mapping)) => mapping,
Ok(None) => {
return JobOutcome::Failed(format!(
"no upstream order is recorded for local order {order_id}"
));
}
Err(error) => {
return JobOutcome::Retry(format!("reading the upstream order failed: {error}"));
}
};
match relay(
inner,
order_id,
&mapping.csr_der,
&mapping.upstream_order_url,
)
.await
{
Ok(chain) => settle(inner, order_id, chain).await,
Err(RelayFailure::Retryable(reason)) => JobOutcome::Retry(reason),
Err(RelayFailure::Permanent(reason)) => JobOutcome::Failed(reason),
}
}
async fn abandon(&self, job: &Job, reason: &str) {
let inner = &self.0;
let Some(order_id) = job.payload.get("order_id").and_then(Value::as_str) else {
return;
};
warn!(event = "upstream_relay_failed", outcome = "failure", order_id = %order_id, reason = %reason);
let mut order = match Order::find_by_id(order_id, &inner.database).await {
Ok(Some(order)) => order,
Ok(None) => {
warn!(event = "upstream_relay_order_vanished", outcome = "failure", order_id = %order_id);
return;
}
Err(error) => {
error!(event = "upstream_relay_order_lookup_failed", outcome = "failure", order_id = %order_id, error = %error);
return;
}
};
let record = relay_record(
crate::audit::AuditEvent::CertificateIssueFailed,
&order,
inner,
)
.await
.with_reason("serverInternal")
.with_detail(reason);
inner.metrics.record_audit(&record);
crate::audit::write(record, &inner.database).await;
let problem = Problem::server_internal("Upstream certificate issuance failed");
if let Err(error) = order
.mark_invalid(problem.to_value(), &inner.database)
.await
{
error!(event = "upstream_relay_mark_invalid_failed", outcome = "failure", order_id = %order_id, error = %error);
}
if let Err(error) = UpstreamOrder::mark_invalid(order_id, reason, &inner.database).await {
warn!(event = "upstream_order_mark_invalid_failed", outcome = "failure", error = %error);
}
}
async fn recover(&self, queue: &JobQueue) {
let inner = &self.0;
let pending = match UpstreamOrder::list_processing(&inner.profiles, &inner.database).await {
Ok(pending) => pending,
Err(error) => {
error!(event = "upstream_resume_lookup_failed", outcome = "failure", error = %error);
return;
}
};
if pending.is_empty() {
return;
}
info!(
event = "upstream_relay_resume_started",
outcome = "progress",
count = pending.len()
);
if pending.len() >= crate::sqlite::upstream_order::MAX_PROCESSING_BATCH {
warn!(
event = "upstream_relay_batch_capped",
outcome = "advisory",
count = pending.len(),
"more orders are still processing than one recovery pass picks up; the \
rest are taken by a later restart",
);
}
for row in pending {
let deadline = order_deadline(&row.order_id, inner).await;
queue
.enqueue_or_log(relay_spec(&row.order_id, deadline))
.await;
}
}
}
pub(super) fn relay_spec(order_id: &str, deadline: Option<i64>) -> JobSpec {
JobSpec::now(RELAY_JOB_KIND, order_id)
.with_payload(json!({ "order_id": order_id }))
.with_deadline(deadline)
}
pub(super) async fn order_deadline(order_id: &str, inner: &Inner) -> Option<i64> {
match Order::find_by_id(order_id, &inner.database).await {
Ok(Some(order)) => Some(order.expires),
Ok(None) => None,
Err(error) => {
warn!(event = "upstream_relay_deadline_lookup_failed", outcome = "failure", order_id = %order_id, error = %error);
None
}
}
}
async fn relay(
inner: &Inner,
order_id: &str,
csr_der: &[u8],
order_url: &str,
) -> Result<String, RelayFailure> {
match &inner.strategy {
ChallengeStrategy::Dns01(updater) => {
let view = poll_until(inner, order_url, &["pending", "ready", "valid"]).await?;
if view.status == "pending" {
answer_dns01(inner, updater.as_ref(), &view.authorizations).await?;
}
}
ChallengeStrategy::Http01(tokens) => {
let view = poll_until(inner, order_url, &["pending", "ready", "valid"]).await?;
if view.status == "pending" {
answer_http01(inner, tokens.clone(), &view.authorizations).await?;
}
}
ChallengeStrategy::Bypass => {
let view = poll_until(inner, order_url, &["pending", "ready", "valid"]).await?;
if view.status == "pending" {
answer_bypass(inner, &view.authorizations).await?;
}
}
}
let view = poll_until(inner, order_url, &["ready", "valid"]).await?;
let view = if view.status == "ready" {
let finalize = view.finalize.clone().ok_or_else(|| {
RelayFailure::Permanent(
"upstream order is ready but advertises no finalize URL".to_string(),
)
})?;
let payload = json!({
"csr": BASE64_URL_SAFE_NO_PAD.encode(csr_der),
});
inner
.client
.post(
&inner.account,
&Signer::Kid(&inner.kid),
&finalize,
Some(&payload),
)
.await
.map_err(|error| classify(&error))?;
poll_until(inner, order_url, &["valid"]).await?
} else {
view
};
let certificate_url = view.certificate.ok_or_else(|| {
RelayFailure::Permanent(
"upstream order is valid but carries no certificate URL".to_string(),
)
})?;
let response = inner
.client
.get(&inner.account, &inner.kid, &certificate_url)
.await
.map_err(|error| classify(&error))?;
let chain = response.text().map_err(|error| classify(&error))?;
if let Err(error) =
UpstreamOrder::mark_valid(order_id, Some(&certificate_url), &inner.database).await
{
warn!(event = "upstream_order_mark_valid_failed", outcome = "failure", error = %error);
}
Ok(chain)
}
fn upstream_thumbprint(inner: &Inner) -> Result<String, RelayFailure> {
crate::extractors::acme::jwk_thumbprint(inner.account.spki_der()).map_err(|error| {
RelayFailure::Permanent(format!(
"cannot derive the upstream account thumbprint: {error}"
))
})
}
async fn answer_dns01(
inner: &Inner,
updater: &dyn dns01::DnsUpdater,
authorizations: &[String],
) -> Result<(), RelayFailure> {
let thumbprint = upstream_thumbprint(inner)?;
for authz_url in authorizations {
let authz = read_authz(inner, authz_url).await?;
if authz.status != "pending" {
continue;
}
let challenge = authz
.challenges
.iter()
.find(|challenge| challenge.typ == crate::challenge::DNS_01)
.ok_or_else(|| {
RelayFailure::Permanent(format!(
"upstream authorization for {} offers no dns-01 challenge",
authz.identifier.value
))
})?;
let name = crate::challenge::dns_01::record_name(&authz.identifier.value);
let key_authorization = format!("{}.{thumbprint}", challenge.token);
let value = crate::challenge::dns_01::expected_value(&key_authorization);
let fqdn = if name.ends_with('.') {
name.clone()
} else {
format!("{name}.")
};
updater.upsert_txt(&fqdn, &value).await.map_err(|error| {
RelayFailure::Retryable(format!("publishing {fqdn} failed: {error}"))
})?;
let triggered = trigger_and_await(inner, &challenge.url, authz_url).await;
if let Err(error) = updater.delete_txt(&fqdn, &value).await {
warn!(event = "signer_relay_dns_01_cleanup_failed", outcome = "failure", name = %fqdn, error = %error);
}
triggered?;
}
Ok(())
}
async fn answer_http01(
inner: &Inner,
tokens: Arc<dyn http01::TokenStore>,
authorizations: &[String],
) -> Result<(), RelayFailure> {
let thumbprint = upstream_thumbprint(inner)?;
for authz_url in authorizations {
let authz = read_authz(inner, authz_url).await?;
if authz.status != "pending" {
continue;
}
if authz.identifier.value.starts_with("*.") {
return Err(RelayFailure::Permanent(format!(
"upstream authorization for {} is a wildcard, which http-01 cannot validate: use \
signer.relay.challenge_strategy = \"dns01\" for wildcard names",
authz.identifier.value
)));
}
let challenge = authz
.challenges
.iter()
.find(|challenge| challenge.typ == crate::challenge::HTTP_01)
.ok_or_else(|| {
RelayFailure::Permanent(format!(
"upstream authorization for {} offers no http-01 challenge",
authz.identifier.value
))
})?;
let key_authorization = format!("{}.{thumbprint}", challenge.token);
let _published =
http01::PublishedToken::publish(tokens.clone(), &challenge.token, &key_authorization);
trigger_and_await(inner, &challenge.url, authz_url).await?;
}
Ok(())
}
async fn answer_bypass(inner: &Inner, authorizations: &[String]) -> Result<(), RelayFailure> {
for authz_url in authorizations {
let authz = read_authz(inner, authz_url).await?;
if authz.status != "pending" {
continue;
}
let challenge = authz.challenges.first().ok_or_else(|| {
RelayFailure::Permanent(format!(
"upstream authorization for {} offers no challenges to bypass",
authz.identifier.value
))
})?;
trigger_and_await(inner, &challenge.url, authz_url).await?;
}
Ok(())
}
async fn read_authz(inner: &Inner, authz_url: &str) -> Result<UpstreamAuthzView, RelayFailure> {
let response = inner
.client
.get(&inner.account, &inner.kid, authz_url)
.await
.map_err(|error| classify(&error))?;
response.json().map_err(|error| classify(&error))
}
async fn trigger_and_await(
inner: &Inner,
challenge_url: &str,
authz_url: &str,
) -> Result<(), RelayFailure> {
inner
.client
.post(
&inner.account,
&Signer::Kid(&inner.kid),
challenge_url,
Some(&json!({})),
)
.await
.map_err(|error| match classify(&error) {
RelayFailure::Retryable(reason) => RelayFailure::Retryable(format!(
"triggering the upstream challenge failed: {reason}"
)),
RelayFailure::Permanent(reason) => RelayFailure::Permanent(format!(
"triggering the upstream challenge failed: {reason}"
)),
})?;
let deadline = tokio::time::Instant::now() + inner.poll.timeout;
loop {
let authz = read_authz(inner, authz_url).await?;
match authz.status.as_str() {
"valid" => return Ok(()),
"invalid" => {
return Err(RelayFailure::Permanent(format!(
"upstream rejected the challenge for {}",
authz.identifier.value
)));
}
_ if tokio::time::Instant::now() >= deadline => {
return Err(RelayFailure::Retryable(format!(
"upstream authorization for {} did not settle in time",
authz.identifier.value
)));
}
_ => tokio::time::sleep(inner.poll.interval).await,
}
}
}
async fn poll_until(
inner: &Inner,
order_url: &str,
wanted: &[&str],
) -> Result<UpstreamOrderView, RelayFailure> {
loop {
let response = inner
.client
.get(&inner.account, &inner.kid, order_url)
.await
.map_err(|error| classify(&error))?;
let view: UpstreamOrderView = response.json().map_err(|error| classify(&error))?;
if wanted.contains(&view.status.as_str()) {
return Ok(view);
}
if view.status == "invalid" {
let detail = view
.error
.as_ref()
.and_then(|error| error.get("detail"))
.and_then(Value::as_str)
.unwrap_or("no detail given")
.to_string();
return Err(RelayFailure::Permanent(format!(
"upstream order became invalid: {detail}"
)));
}
let wait = response
.retry_after
.map(Duration::from_secs)
.unwrap_or(inner.poll.interval);
tokio::time::sleep(wait).await;
}
}
pub(super) async fn settle(inner: &Inner, order_id: &str, chain: String) -> JobOutcome {
let mut order = match Order::find_by_id(order_id, &inner.database).await {
Ok(Some(order)) => order,
Ok(None) => {
warn!(event = "upstream_relay_order_vanished", outcome = "failure", order_id = %order_id);
return JobOutcome::Failed("the local order no longer exists".to_string());
}
Err(error) => {
error!(event = "upstream_relay_order_lookup_failed", outcome = "failure", order_id = %order_id, error = %error);
return JobOutcome::Retry(format!("reading the local order failed: {error}"));
}
};
let leaf = match crate::cert::leaf_der_from_chain(&chain) {
Ok(leaf) => leaf,
Err(error) => {
return JobOutcome::Failed(format!("upstream chain unparsable: {error}"));
}
};
let (serial, pubkey) = match crate::cert::cert_serial_and_spki(&leaf) {
Ok(parts) => parts,
Err(error) => {
return JobOutcome::Failed(format!("upstream leaf unparsable: {error}"));
}
};
if let Err(error) = order
.finalize(chain, serial.clone(), pubkey, &inner.database)
.await
{
error!(event = "upstream_relay_finalize_failed", outcome = "failure", order_id = %order_id, error = %error);
return JobOutcome::Retry(format!("recording the certificate failed: {error}"));
}
info!(event = "upstream_relay_succeeded", outcome = "success", order_id = %order_id, cert_serial = %serial);
let record = relay_record(crate::audit::AuditEvent::CertificateIssued, &order, inner)
.await
.with_serial(serial.clone());
inner.metrics.record_audit(&record);
crate::audit::write(record, &inner.database).await;
if let Some(dispatcher) = inner.notifiers.get(&order.profile) {
dispatcher
.dispatch(NotifyEvent::CertificateIssued(CertificateIssuedData {
profile: order.profile.clone(),
order_id: order_id.to_string(),
account_id: order.account_id.clone(),
cert_serial: serial.clone(),
identifiers: order.identifiers.iter().map(|i| i.value.clone()).collect(),
client_ip: None,
}))
.await;
}
JobOutcome::Done
}
async fn relay_record(
event: crate::audit::AuditEvent,
order: &Order,
inner: &Inner,
) -> crate::audit::AuditRecord {
let mapping = UpstreamOrder::find_by_order_id(&order.id, &inner.database)
.await
.unwrap_or_else(|error| {
warn!(event = "upstream_order_client_context_lookup_failed", outcome = "failure", order_id = %order.id, error = %error);
None
});
let (actor, client) = match &mapping {
Some(mapping) => (
crate::audit::Actor::acme(&order.account_id),
mapping.client(),
),
None => (
crate::audit::Actor::system(),
crate::audit::ClientContext::default(),
),
};
crate::audit::AuditRecord::new(event, &order.profile, actor)
.with_order(order)
.with_client(client)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn every_upstream_failure_is_classified_by_whose_answer_it_is() {
let permanent = |error: UpstreamError| match classify(&error) {
RelayFailure::Permanent(_) => (),
other => panic!("{error} must be permanent, got {other:?}"),
};
let retryable = |error: UpstreamError| match classify(&error) {
RelayFailure::Retryable(_) => (),
other => panic!("{error} must be retryable, got {other:?}"),
};
retryable(UpstreamError::Transport("connection reset".to_string()));
retryable(UpstreamError::Url("no host".to_string()));
retryable(UpstreamError::Protocol(
"expected a JSON object".to_string(),
));
retryable(UpstreamError::Problem {
status: 500,
typ: "urn:ietf:params:acme:error:serverInternal".to_string(),
detail: "internal error".to_string(),
});
retryable(UpstreamError::Problem {
status: 503,
typ: "urn:ietf:params:acme:error:serverInternal".to_string(),
detail: "try later".to_string(),
});
retryable(UpstreamError::Problem {
status: 429,
typ: "urn:ietf:params:acme:error:rateLimited".to_string(),
detail: "too many certificates".to_string(),
});
permanent(UpstreamError::Problem {
status: 403,
typ: "urn:ietf:params:acme:error:unauthorized".to_string(),
detail: "not authorized".to_string(),
});
permanent(UpstreamError::Problem {
status: 400,
typ: "urn:ietf:params:acme:error:badCSR".to_string(),
detail: "unacceptable key".to_string(),
});
permanent(UpstreamError::Problem {
status: 400,
typ: "urn:ietf:params:acme:error:rejectedIdentifier".to_string(),
detail: "will not issue for that name".to_string(),
});
permanent(UpstreamError::Jws("signing failed".to_string()));
}
#[test]
fn a_classified_failure_keeps_the_upstream_error_text() {
let error = UpstreamError::Transport("connection reset".to_string());
let failure = classify(&error);
assert_eq!(failure.reason(), error.to_string());
assert_eq!(failure.to_string(), error.to_string());
}
#[test]
fn a_relay_job_spec_carries_the_order_id_as_both_identity_and_payload() {
let spec = relay_spec("ord-1", Some(1_234));
assert_eq!(spec.kind, RELAY_JOB_KIND);
assert_eq!(spec.key, "ord-1");
assert_eq!(spec.payload, json!({"order_id": "ord-1"}));
assert_eq!(spec.deadline, Some(1_234));
}
}