use std::net::IpAddr;
use std::sync::Arc;
use std::time::Duration;
use tracing::{error, warn};
use super::order::OrderService;
use crate::profile::Profile;
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_store::account::Account;
use acme_proxy_store::authz::Authorization;
use acme_proxy_store::authz::Challenge;
use acme_proxy_store::db::Database;
use acme_proxy_store::job::Job;
use acme_proxy_store::nonce::now_secs;
use acme_proxy_store::order::Order;
use acme_proxy_store::status::ChallengeStatus;
pub const CHALLENGE_VALIDATE_KIND: &str = "challenge_validate";
#[must_use]
pub fn challenge_validate_spec(
challenge_id: &str,
client_ip: Option<IpAddr>,
authz_expires: i64,
) -> JobSpec {
JobSpec::now(CHALLENGE_VALIDATE_KIND, challenge_id)
.with_payload(serde_json::json!({
"challenge_id": challenge_id,
"client_ip": client_ip.map(|ip| acme_proxy_core::client::canonical(ip).to_string()),
}))
.with_deadline(Some(authz_expires))
}
pub struct ChallengeValidateJob {
database: Arc<Database>,
audit: Arc<Auditor>,
profiles: Vec<(String, Arc<Profile>)>,
}
struct Subject {
challenge: Challenge,
authz: Authorization,
order: Order,
account: Account,
profile: Arc<Profile>,
}
impl ChallengeValidateJob {
#[must_use]
pub fn new(
database: Arc<Database>,
audit: Arc<Auditor>,
profiles: Vec<(String, Arc<Profile>)>,
) -> Self {
Self {
database,
audit,
profiles,
}
}
fn profile(&self, name: &str) -> Option<&Arc<Profile>> {
super::mounted(&self.profiles, name)
}
async fn load(&self, challenge_id: &str) -> Result<Subject, JobOutcome> {
let challenge = match Challenge::find_by_id(challenge_id, &self.database).await {
Ok(Some(challenge)) => challenge,
Ok(None) => {
return Err(JobOutcome::Failed(
"the challenge no longer exists".to_string(),
));
}
Err(error) => {
return Err(JobOutcome::Retry(format!(
"reading the challenge failed: {error}"
)));
}
};
let authz = match Authorization::find_by_id(
challenge.authz_id.to_string().as_str(),
&self.database,
)
.await
{
Ok(Some(authz)) => authz,
Ok(None) => {
return Err(JobOutcome::Failed(
"the authorization no longer exists".to_string(),
));
}
Err(error) => {
return Err(JobOutcome::Retry(format!(
"reading the authorization failed: {error}"
)));
}
};
let order =
match Order::find_by_id(authz.order_id.to_string().as_str(), &self.database).await {
Ok(Some(order)) => order,
Ok(None) => {
return Err(JobOutcome::Failed("the order no longer exists".to_string()));
}
Err(error) => {
return Err(JobOutcome::Retry(format!(
"reading the order failed: {error}"
)));
}
};
let Some(profile) = self.profile(&order.profile) else {
return Err(JobOutcome::Retry(format!(
"profile `{}` is not mounted by this process",
order.profile
)));
};
let account = match Account::find_by_id(
&order.profile,
order.account_id.to_string().as_str(),
&self.database,
)
.await
{
Ok(Some(account)) => account,
Ok(None) => {
return Err(JobOutcome::Failed(
"the account no longer exists".to_string(),
));
}
Err(error) => {
return Err(JobOutcome::Retry(format!(
"reading the account failed: {error}"
)));
}
};
Ok(Subject {
challenge,
authz,
order,
account,
profile: profile.clone(),
})
}
}
#[async_trait::async_trait]
impl JobHandler for ChallengeValidateJob {
fn kind(&self) -> &'static str {
CHALLENGE_VALIDATE_KIND
}
async fn run(&self, job: &Job) -> JobOutcome {
let Some(challenge_id) = job.payload["challenge_id"].as_str() else {
return JobOutcome::Failed("the payload names no challenge".to_string());
};
let client_ip = job.payload["client_ip"]
.as_str()
.and_then(|ip| ip.parse::<IpAddr>().ok());
let Subject {
mut challenge,
mut authz,
mut order,
account,
profile,
} = match self.load(challenge_id).await {
Ok(subject) => subject,
Err(outcome) => return outcome,
};
if challenge.status != ChallengeStatus::Processing {
return JobOutcome::Done;
}
let orders = OrderService {
database: &self.database,
audit: &self.audit,
profile: &profile,
};
match orders
.run_validation(&account, &mut challenge, &mut authz, &mut order, client_ip)
.await
{
Ok(()) => JobOutcome::Done,
Err(_) => JobOutcome::Retry(
"the validation answer could not be computed or stored".to_string(),
),
}
}
fn lease(&self, _job: &Job) -> Option<Duration> {
self.profiles
.iter()
.map(|(_, profile)| profile.challenges.timeout())
.max()
.map(|timeout| timeout + LEASE_HEADROOM)
}
async fn abandon(&self, job: &Job, reason: &str) {
let Some(challenge_id) = job.payload["challenge_id"].as_str() else {
return;
};
warn!(
event = "challenge_validation_abandoned",
outcome = "failure",
challenge_id = %challenge_id,
attempts = job.attempts,
reason = %reason,
"a queued challenge validation was given up; its order is marked invalid"
);
let Ok(Subject {
mut challenge,
mut authz,
mut order,
profile,
..
}) = self.load(challenge_id).await
else {
return;
};
if challenge.status != ChallengeStatus::Processing {
return;
}
let orders = OrderService {
database: &self.database,
audit: &self.audit,
profile: &profile,
};
if orders
.abandon_validation(&mut challenge, &mut authz, &mut order, reason)
.await
.is_err()
{
error!(
event = "challenge_validation_abandon_failed",
outcome = "failure",
challenge_id = %challenge_id,
"the challenge could not be marked invalid and will stay `processing` \
until its authorization expires"
);
}
}
async fn recover(&self, queue: &JobQueue) {
let stranded = match Challenge::find_processing(now_secs(), &self.database).await {
Ok(rows) => rows,
Err(error) => {
error!(
event = "challenge_validation_recovery_failed",
outcome = "failure",
error = %error
);
return;
}
};
for (challenge_id, authz_expires) in stranded {
queue
.enqueue_or_log(challenge_validate_spec(
challenge_id.to_string().as_str(),
None,
authz_expires,
))
.await;
}
}
}
const LEASE_HEADROOM: Duration = Duration::from_secs(30);
#[cfg(test)]
mod tests {
use super::*;
use crate::acme::order::tests::{account, profile};
use acme_proxy_core::identifier::Identifier;
use acme_proxy_net::challenge::ChallengeError;
use acme_proxy_net::challenge::ChallengeRegistry;
use acme_proxy_net::challenge::ChallengeValidator;
use acme_proxy_net::challenge::ValidationContext;
use acme_proxy_store::status::AuthzStatus;
use acme_proxy_store::status::OrderStatus;
struct Refusing;
#[async_trait::async_trait]
impl ChallengeValidator for Refusing {
fn typ(&self) -> &'static str {
"http-01"
}
async fn validate(&self, _ctx: &ValidationContext<'_>) -> Result<(), ChallengeError> {
Err(ChallengeError::IncorrectResponse("wrong body".into()))
}
}
async fn subject(
database: &Arc<Database>,
account: &Account,
) -> (Order, Authorization, Challenge) {
let order = Order::create(
"default",
account.id,
acme_proxy_store::testutil::dns_identifiers(&["a.example.com"]),
now_secs() + 3600,
None,
None,
database,
)
.await
.unwrap();
let authz = Authorization::create(
order.id,
Identifier::dns("a.example.com"),
order.expires,
database,
)
.await
.unwrap();
let challenge = Challenge::create(authz.id, "http-01", database)
.await
.unwrap();
(order, authz, challenge)
}
fn handler(database: &Arc<Database>, challenges: ChallengeRegistry) -> ChallengeValidateJob {
let profile = Arc::new(profile(database, challenges));
ChallengeValidateJob::new(
database.clone(),
Arc::new(Auditor::offline(database.clone())),
vec![(profile.name.clone(), profile)],
)
}
fn row(challenge_id: &str) -> Job {
let spec = challenge_validate_spec(challenge_id, None, now_secs() + 3600);
Job {
dedup_key: spec.key.clone(),
payload: spec.payload.clone(),
..acme_proxy_store::testutil::job_fixture()
}
}
#[tokio::test]
async fn a_claimed_challenge_is_validated_and_its_order_promoted() {
let database = Arc::new(Database::connect_in_memory().await.unwrap());
let account = account(&database).await;
let (order, authz, mut challenge) = subject(&database, &account).await;
assert_eq!(
challenge.claim_for_validation(0, &database).await.unwrap(),
acme_proxy_store::authz::ValidationClaim::Claimed
);
let job = handler(&database, ChallengeRegistry::default());
assert!(matches!(
job.run(&row(challenge.id.to_string().as_str())).await,
JobOutcome::Done
));
let reloaded = Challenge::find_by_id(challenge.id.to_string().as_str(), &database)
.await
.unwrap()
.unwrap();
assert_eq!(reloaded.status, ChallengeStatus::Valid);
assert_eq!(
Authorization::find_by_id(authz.id.to_string().as_str(), &database)
.await
.unwrap()
.unwrap()
.status,
AuthzStatus::Valid
);
assert_eq!(
Order::find_by_id(order.id.to_string().as_str(), &database)
.await
.unwrap()
.unwrap()
.status,
OrderStatus::Ready
);
}
#[tokio::test]
async fn a_refused_validation_is_recorded_and_the_job_is_done() {
let database = Arc::new(Database::connect_in_memory().await.unwrap());
let account = account(&database).await;
let (order, _, mut challenge) = subject(&database, &account).await;
assert_eq!(
challenge.claim_for_validation(0, &database).await.unwrap(),
acme_proxy_store::authz::ValidationClaim::Claimed
);
let registry = ChallengeRegistry::new(
vec![Arc::new(Refusing)],
vec!["http-01".to_string()],
false,
Duration::from_secs(5),
);
let job = handler(&database, registry);
assert!(matches!(
job.run(&row(challenge.id.to_string().as_str())).await,
JobOutcome::Done
));
assert_eq!(
Order::find_by_id(order.id.to_string().as_str(), &database)
.await
.unwrap()
.unwrap()
.status,
OrderStatus::Invalid
);
}
#[tokio::test]
async fn an_unmounted_profile_is_retried() {
let database = Arc::new(Database::connect_in_memory().await.unwrap());
let account = account(&database).await;
let (_, _, mut challenge) = subject(&database, &account).await;
assert_eq!(
challenge.claim_for_validation(0, &database).await.unwrap(),
acme_proxy_store::authz::ValidationClaim::Claimed
);
let job = ChallengeValidateJob::new(
database.clone(),
Arc::new(Auditor::offline(database.clone())),
vec![],
);
let JobOutcome::Retry(reason) = job.run(&row(challenge.id.to_string().as_str())).await
else {
panic!("an unmounted profile must be retried, never failed");
};
assert!(reason.contains("not mounted"), "{reason}");
}
#[tokio::test]
async fn a_vanished_challenge_is_failed() {
let database = Arc::new(Database::connect_in_memory().await.unwrap());
let job = handler(&database, ChallengeRegistry::default());
let missing = uuid::Uuid::now_v7().to_string();
assert!(matches!(
job.run(&row(&missing)).await,
JobOutcome::Failed(_)
));
}
#[tokio::test]
async fn a_vanished_parent_is_failed_and_an_unreadable_one_retried() {
for (table, reason, delete, drop) in [
(
"authorizations",
"authorization",
"DELETE FROM authorizations;",
"DROP TABLE authorizations;",
),
(
"orders",
"order",
"DELETE FROM orders;",
"DROP TABLE orders;",
),
(
"accounts",
"account",
"DELETE FROM accounts;",
"DROP TABLE accounts;",
),
] {
for drop_table in [false, true] {
let database = Arc::new(Database::connect_in_memory().await.unwrap());
let account = account(&database).await;
let (_, _, challenge) = subject(&database, &account).await;
let target = if drop_table { drop } else { delete };
for statement in ["PRAGMA foreign_keys = OFF;", target] {
sqlx::query(statement)
.execute(database.raw_pool())
.await
.unwrap();
}
let job = handler(&database, ChallengeRegistry::default());
let outcome = job.run(&row(challenge.id.to_string().as_str())).await;
match (drop_table, outcome) {
(false, JobOutcome::Failed(why)) => {
assert!(why.contains(reason), "{table}: {why}");
}
(true, JobOutcome::Retry(why)) => {
assert!(why.contains(reason), "{table}: {why}");
}
(_, other) => panic!("{table}, dropped {drop_table}: {other:?}"),
}
}
}
}
#[tokio::test]
async fn a_payload_naming_no_challenge_is_failed() {
let database = Arc::new(Database::connect_in_memory().await.unwrap());
let job = handler(&database, ChallengeRegistry::default());
let row = Job {
payload: serde_json::json!({}),
..acme_proxy_store::testutil::job_fixture()
};
assert!(matches!(job.run(&row).await, JobOutcome::Failed(_)));
}
#[tokio::test]
async fn a_settled_challenge_is_done_without_revalidating() {
let database = Arc::new(Database::connect_in_memory().await.unwrap());
let account = account(&database).await;
let (_, _, mut challenge) = subject(&database, &account).await;
assert_eq!(
challenge.claim_for_validation(0, &database).await.unwrap(),
acme_proxy_store::authz::ValidationClaim::Claimed
);
let job = handler(&database, ChallengeRegistry::default());
assert!(matches!(
job.run(&row(challenge.id.to_string().as_str())).await,
JobOutcome::Done
));
let registry = ChallengeRegistry::new(
vec![Arc::new(Refusing)],
vec!["http-01".to_string()],
false,
Duration::from_secs(5),
);
let job = handler(&database, registry);
assert!(matches!(
job.run(&row(challenge.id.to_string().as_str())).await,
JobOutcome::Done
));
assert_eq!(
Challenge::find_by_id(challenge.id.to_string().as_str(), &database)
.await
.unwrap()
.unwrap()
.status,
ChallengeStatus::Valid
);
}
#[tokio::test]
async fn abandoning_marks_the_challenge_and_its_order_invalid() {
let database = Arc::new(Database::connect_in_memory().await.unwrap());
let account = account(&database).await;
let (order, authz, mut challenge) = subject(&database, &account).await;
assert_eq!(
challenge.claim_for_validation(0, &database).await.unwrap(),
acme_proxy_store::authz::ValidationClaim::Claimed
);
let job = handler(&database, ChallengeRegistry::default());
job.abandon(
&row(challenge.id.to_string().as_str()),
"the attempts ran out",
)
.await;
let reloaded = Challenge::find_by_id(challenge.id.to_string().as_str(), &database)
.await
.unwrap()
.unwrap();
assert_eq!(reloaded.status, ChallengeStatus::Invalid);
assert_eq!(
Authorization::find_by_id(authz.id.to_string().as_str(), &database)
.await
.unwrap()
.unwrap()
.status,
AuthzStatus::Invalid
);
assert_eq!(
Order::find_by_id(order.id.to_string().as_str(), &database)
.await
.unwrap()
.unwrap()
.status,
OrderStatus::Invalid
);
}
#[tokio::test]
async fn recover_requeues_a_stranded_claim_exactly_once() {
let database = Arc::new(Database::connect_in_memory().await.unwrap());
let account = account(&database).await;
let (_, _, mut challenge) = subject(&database, &account).await;
assert_eq!(
challenge.claim_for_validation(0, &database).await.unwrap(),
acme_proxy_store::authz::ValidationClaim::Claimed
);
let queue = acme_proxy_jobs::testutil::idle_job_queue(database.clone());
let job = handler(&database, ChallengeRegistry::default());
job.recover(&queue).await;
assert_eq!(
acme_proxy_store::job::Job::count_live(CHALLENGE_VALIDATE_KIND, &database)
.await
.unwrap(),
1
);
job.recover(&queue).await;
assert_eq!(
acme_proxy_store::job::Job::count_live(CHALLENGE_VALIDATE_KIND, &database)
.await
.unwrap(),
1
);
}
#[tokio::test]
async fn recover_ignores_a_challenge_nobody_claimed() {
let database = Arc::new(Database::connect_in_memory().await.unwrap());
let account = account(&database).await;
let (_, _, _challenge) = subject(&database, &account).await;
let queue = acme_proxy_jobs::testutil::idle_job_queue(database.clone());
handler(&database, ChallengeRegistry::default())
.recover(&queue)
.await;
assert_eq!(
acme_proxy_store::job::Job::count_live(CHALLENGE_VALIDATE_KIND, &database)
.await
.unwrap(),
0
);
}
}