use std::sync::Arc;
use std::time::Duration;
use async_trait::async_trait;
use base64::prelude::*;
use serde_json::json;
use tracing::{debug, info, warn};
use crate::config::RelayConfig;
use crate::jobs::JobQueue;
use crate::signer::{IssueOutcome, RenewalWindow, RequestedValidity, SignerBackend, SignerError};
use crate::sqlite::db::Database;
use crate::sqlite::order::Identifier;
use crate::sqlite::upstream_order::UpstreamOrder;
pub mod account;
pub mod client;
pub mod dns01;
pub mod eab;
pub mod flow;
pub mod http01;
#[cfg(test)]
pub mod testsrv;
pub mod wire;
use client::{AccountKey, AcmeClient, Signer};
use account::provision;
pub use account::{register_upstream_account, stored_kid};
pub(crate) use eab::decode_secret;
use flow::{OrderContext, relay_spec};
pub(crate) use flow::{RELAY_JOB_KIND, abandon_relayed_order};
use wire::{RenewalInfoView, UpstreamOrderView, parse_rfc3339, upstream_to_signer_error};
pub enum ChallengeStrategy {
Bypass,
Dns01(Arc<dyn dns01::DnsUpdater>),
Http01(Arc<dyn http01::TokenStore>),
}
fn token_store_key(account_key_path: &str) -> String {
format!("relay.http01:{account_key_path}")
}
struct PollConfig {
interval: Duration,
timeout: Duration,
}
struct Inner {
client: AcmeClient,
account: AccountKey,
kid: String,
account_key_path: String,
http01_tokens: Option<Arc<http01::MemoryTokenStore>>,
database: Arc<Database>,
strategy: ChallengeStrategy,
poll: PollConfig,
notifiers: crate::notify::Notifiers,
metrics: Arc<crate::metrics::Metrics>,
jobs: JobQueue,
}
pub struct RelaySigner(Arc<Inner>);
pub struct RelayState(Arc<Inner>);
impl RelaySigner {
pub fn from_config(
cfg: &RelayConfig,
parts: &crate::signer::SignerParts,
carried: &crate::signer::CarriedState,
) -> anyhow::Result<Self> {
let outbound = parts.egress.outbound();
if cfg.directory_url.is_empty() {
anyhow::bail!(
"signer.relay.directory_url is empty: the relay backend has no upstream \
to relay to"
);
}
let poll = PollConfig {
interval: Duration::from_millis(cfg.poll_interval_ms),
timeout: Duration::from_secs(cfg.poll_timeout_secs),
};
let (client, account, kid, strategy, http01_tokens) = std::thread::scope(|scope| {
scope
.spawn(|| -> anyhow::Result<_> {
let (strategy, http01_tokens) = match cfg.challenge_strategy.as_str() {
"bypass" => (ChallengeStrategy::Bypass, None),
"dns01" => match cfg.dns01.provider.as_str() {
"rfc2136" => (
ChallengeStrategy::Dns01(Arc::new(
dns01::Rfc2136Updater::from_config(&cfg.dns01.rfc2136)?,
)),
None,
),
other => anyhow::bail!(
"unknown signer.relay.dns01.provider: {other} (supported: rfc2136)"
),
},
"http01" => {
info!(
event = "signer_relay_http_01_selected",
outcome = "advisory",
path = crate::challenge::http_01::WELL_KNOWN_PREFIX,
"the upstream will fetch \
http://<identifier>:80/.well-known/acme-challenge/<token>; a \
reverse proxy must forward or redirect that path to this server \
(RFC 8555 §8.3 permits a redirect, so it need not share the name)"
);
let key = token_store_key(&cfg.account_key_path);
let tokens = match carried.get::<http01::MemoryTokenStore>(&key) {
Some(tokens) => {
info!(
event = "signer_relay_token_store_adopted",
outcome = "success",
account_key_path = %cfg.account_key_path,
"key authorizations published for the upstream survive \
the reload that rebuilt this backend"
);
tokens
}
None => Arc::new(http01::MemoryTokenStore::new()),
};
(ChallengeStrategy::Http01(tokens.clone()), Some(tokens))
}
other => anyhow::bail!(
"unknown signer.relay.challenge_strategy: {other} \
(supported: bypass, dns01, http01)"
),
};
let (client, account, kid) = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()?
.block_on(provision(cfg, outbound, poll.timeout))?;
Ok((client, account, kid, strategy, http01_tokens))
})
.join()
.unwrap_or_else(|_| Err(anyhow::anyhow!("upstream provisioning thread panicked")))
})?;
Ok(Self(Arc::new(Inner {
client,
account,
kid,
account_key_path: cfg.account_key_path.clone(),
http01_tokens,
database: parts.database.clone(),
strategy,
poll,
notifiers: parts.notifiers.clone(),
metrics: parts.metrics.clone(),
jobs: parts.jobs.clone(),
})))
}
}
#[async_trait]
impl SignerBackend for RelaySigner {
#[tracing::instrument(name = "relay_issue", skip_all, fields(order_id = %order_id))]
async fn issue(
&self,
order_id: &str,
csr_der: &[u8],
identifiers: &[Identifier],
validity: RequestedValidity,
) -> Result<IssueOutcome, SignerError> {
let _ = validity;
let inner = self.0.clone();
let payload = json!({
"identifiers": identifiers.iter().map(|identifier| json!({
"type": identifier.typ,
"value": identifier.value,
})).collect::<Vec<_>>(),
});
let response = inner
.client
.post(
&inner.account,
&Signer::Kid(&inner.kid),
&inner.client.directory().new_order.clone(),
Some(&payload),
)
.await
.map_err(upstream_to_signer_error)?;
let order_url = response.location.clone().ok_or_else(|| {
SignerError::Internal("upstream newOrder returned no Location header".to_string())
})?;
let view: UpstreamOrderView = response.json().map_err(upstream_to_signer_error)?;
let created = UpstreamOrder::create(
order_id,
&order_url,
view.finalize.as_deref(),
csr_der,
&inner.database,
)
.await
.map_err(|error| SignerError::Internal(format!("recording upstream order: {error}")))?;
if created.is_none() {
warn!(event = "upstream_relay_already_in_flight", outcome = "advisory", order_id = %order_id);
return Ok(IssueOutcome::Processing);
}
info!(event = "upstream_order_opened", outcome = "success", order_id = %order_id, upstream_url = %order_url);
let context = OrderContext::read(order_id, &inner).await;
inner
.jobs
.enqueue(relay_spec(order_id, &context))
.await
.map_err(|error| SignerError::Internal(format!("queueing the relay: {error}")))?;
Ok(IssueOutcome::Processing)
}
fn relay_state(&self) -> Option<RelayState> {
Some(RelayState(self.0.clone()))
}
#[tracing::instrument(name = "relay_revoke", skip_all)]
async fn revoke(&self, cert_der: &[u8], reason: Option<u32>) -> Result<(), SignerError> {
let inner = &self.0;
let revoke_url = inner
.client
.directory()
.revoke_cert
.clone()
.ok_or_else(|| {
SignerError::Internal("upstream directory advertises no revokeCert".to_string())
})?;
let mut payload = json!({
"certificate": BASE64_URL_SAFE_NO_PAD.encode(cert_der),
});
if let Some(reason) = reason {
payload["reason"] = json!(reason);
}
match inner
.client
.post(
&inner.account,
&Signer::Kid(&inner.kid),
&revoke_url,
Some(&payload),
)
.await
{
Ok(_) => Ok(()),
Err(error) if error.is_already_revoked() => {
debug!(event = "upstream_already_revoked", outcome = "success");
Ok(())
}
Err(error) => Err(upstream_to_signer_error(error)),
}
}
#[tracing::instrument(name = "relay_renewal_info", skip_all)]
async fn renewal_info(&self, cert_der: &[u8]) -> Result<Option<RenewalWindow>, SignerError> {
let inner = &self.0;
let Some(base) = inner.client.directory().renewal_info.clone() else {
debug!(event = "upstream_has_no_renewal_info", outcome = "success");
return Ok(None);
};
let cert_id = match crate::cert::ari_cert_id(cert_der) {
Ok(cert_id) => cert_id,
Err(error) => {
debug!(event = "upstream_renewal_info_cert_id_underivable", outcome = "failure", error = %error);
return Ok(None);
}
};
let url = format!("{}/{cert_id}", base.trim_end_matches('/'));
let response = inner
.client
.get_unsigned(&url)
.await
.map_err(upstream_to_signer_error)?;
let info: RenewalInfoView = response.json().map_err(upstream_to_signer_error)?;
let start = parse_rfc3339(&info.suggested_window.start)?;
let end = parse_rfc3339(&info.suggested_window.end)?;
info!(
event = "upstream_renewal_info_used",
outcome = "success",
start,
end,
explanation_url = ?info.explanation_url,
);
Ok(Some(RenewalWindow {
start,
end,
explanation_url: info.explanation_url,
}))
}
fn http01_tokens(&self) -> Option<Arc<dyn crate::signer::Http01TokenStore>> {
match &self.0.strategy {
ChallengeStrategy::Http01(tokens) => Some(tokens.clone()),
ChallengeStrategy::Bypass | ChallengeStrategy::Dns01(_) => None,
}
}
fn carried_state(&self) -> crate::signer::CarriedState {
let mut carried = crate::signer::CarriedState::new();
if let Some(tokens) = &self.0.http01_tokens {
carried.insert(token_store_key(&self.0.account_key_path), tokens.clone());
}
carried
}
}
#[cfg(test)]
mod tests;