1use std::net::IpAddr;
21use std::sync::Arc;
22use std::time::Duration;
23
24use tracing::{error, warn};
25
26use super::order::OrderService;
27use crate::profile::Profile;
28use acme_proxy_jobs::auditor::Auditor;
29use acme_proxy_jobs::jobs::JobHandler;
30use acme_proxy_jobs::jobs::JobOutcome;
31use acme_proxy_jobs::jobs::JobQueue;
32use acme_proxy_jobs::jobs::JobSpec;
33use acme_proxy_store::account::Account;
34use acme_proxy_store::authz::Authorization;
35use acme_proxy_store::authz::Challenge;
36use acme_proxy_store::db::Database;
37use acme_proxy_store::job::Job;
38use acme_proxy_store::nonce::now_secs;
39use acme_proxy_store::order::Order;
40use acme_proxy_store::status::ChallengeStatus;
41
42pub const CHALLENGE_VALIDATE_KIND: &str = "challenge_validate";
44
45#[must_use]
61pub fn challenge_validate_spec(
62 challenge_id: &str,
63 client_ip: Option<IpAddr>,
64 authz_expires: i64,
65) -> JobSpec {
66 JobSpec::now(CHALLENGE_VALIDATE_KIND, challenge_id)
67 .with_payload(serde_json::json!({
68 "challenge_id": challenge_id,
69 "client_ip": client_ip.map(|ip| acme_proxy_core::client::canonical(ip).to_string()),
70 }))
71 .with_deadline(Some(authz_expires))
72}
73
74pub struct ChallengeValidateJob {
76 database: Arc<Database>,
77 audit: Arc<Auditor>,
78 profiles: Vec<(String, Arc<Profile>)>,
79}
80
81struct Subject {
83 challenge: Challenge,
84 authz: Authorization,
85 order: Order,
86 account: Account,
87 profile: Arc<Profile>,
88}
89
90impl ChallengeValidateJob {
91 #[must_use]
93 pub fn new(
94 database: Arc<Database>,
95 audit: Arc<Auditor>,
96 profiles: Vec<(String, Arc<Profile>)>,
97 ) -> Self {
98 Self {
99 database,
100 audit,
101 profiles,
102 }
103 }
104
105 fn profile(&self, name: &str) -> Option<&Arc<Profile>> {
107 super::mounted(&self.profiles, name)
108 }
109
110 async fn load(&self, challenge_id: &str) -> Result<Subject, JobOutcome> {
115 let challenge = match Challenge::find_by_id(challenge_id, &self.database).await {
116 Ok(Some(challenge)) => challenge,
117 Ok(None) => {
118 return Err(JobOutcome::Failed(
119 "the challenge no longer exists".to_string(),
120 ));
121 }
122 Err(error) => {
123 return Err(JobOutcome::Retry(format!(
124 "reading the challenge failed: {error}"
125 )));
126 }
127 };
128
129 let authz = match Authorization::find_by_id(
130 challenge.authz_id.to_string().as_str(),
131 &self.database,
132 )
133 .await
134 {
135 Ok(Some(authz)) => authz,
136 Ok(None) => {
137 return Err(JobOutcome::Failed(
138 "the authorization no longer exists".to_string(),
139 ));
140 }
141 Err(error) => {
142 return Err(JobOutcome::Retry(format!(
143 "reading the authorization failed: {error}"
144 )));
145 }
146 };
147
148 let order =
149 match Order::find_by_id(authz.order_id.to_string().as_str(), &self.database).await {
150 Ok(Some(order)) => order,
151 Ok(None) => {
152 return Err(JobOutcome::Failed("the order no longer exists".to_string()));
153 }
154 Err(error) => {
155 return Err(JobOutcome::Retry(format!(
156 "reading the order failed: {error}"
157 )));
158 }
159 };
160
161 let Some(profile) = self.profile(&order.profile) else {
164 return Err(JobOutcome::Retry(format!(
165 "profile `{}` is not mounted by this process",
166 order.profile
167 )));
168 };
169
170 let account = match Account::find_by_id(
171 &order.profile,
172 order.account_id.to_string().as_str(),
173 &self.database,
174 )
175 .await
176 {
177 Ok(Some(account)) => account,
178 Ok(None) => {
179 return Err(JobOutcome::Failed(
180 "the account no longer exists".to_string(),
181 ));
182 }
183 Err(error) => {
184 return Err(JobOutcome::Retry(format!(
185 "reading the account failed: {error}"
186 )));
187 }
188 };
189
190 Ok(Subject {
191 challenge,
192 authz,
193 order,
194 account,
195 profile: profile.clone(),
196 })
197 }
198}
199
200#[async_trait::async_trait]
201impl JobHandler for ChallengeValidateJob {
202 fn kind(&self) -> &'static str {
203 CHALLENGE_VALIDATE_KIND
204 }
205
206 async fn run(&self, job: &Job) -> JobOutcome {
217 let Some(challenge_id) = job.payload["challenge_id"].as_str() else {
218 return JobOutcome::Failed("the payload names no challenge".to_string());
219 };
220 let client_ip = job.payload["client_ip"]
221 .as_str()
222 .and_then(|ip| ip.parse::<IpAddr>().ok());
223
224 let Subject {
225 mut challenge,
226 mut authz,
227 mut order,
228 account,
229 profile,
230 } = match self.load(challenge_id).await {
231 Ok(subject) => subject,
232 Err(outcome) => return outcome,
233 };
234
235 if challenge.status != ChallengeStatus::Processing {
239 return JobOutcome::Done;
240 }
241
242 let orders = OrderService {
243 database: &self.database,
244 audit: &self.audit,
245 profile: &profile,
246 };
247 match orders
248 .run_validation(&account, &mut challenge, &mut authz, &mut order, client_ip)
249 .await
250 {
251 Ok(()) => JobOutcome::Done,
252 Err(_) => JobOutcome::Retry(
255 "the validation answer could not be computed or stored".to_string(),
256 ),
257 }
258 }
259
260 fn lease(&self, _job: &Job) -> Option<Duration> {
269 self.profiles
270 .iter()
271 .map(|(_, profile)| profile.challenges.timeout())
272 .max()
273 .map(|timeout| timeout + LEASE_HEADROOM)
274 }
275
276 async fn abandon(&self, job: &Job, reason: &str) {
283 let Some(challenge_id) = job.payload["challenge_id"].as_str() else {
284 return;
285 };
286 warn!(
287 event = "challenge_validation_abandoned",
288 outcome = "failure",
289 challenge_id = %challenge_id,
290 attempts = job.attempts,
291 reason = %reason,
292 "a queued challenge validation was given up; its order is marked invalid"
293 );
294
295 let Ok(Subject {
296 mut challenge,
297 mut authz,
298 mut order,
299 profile,
300 ..
301 }) = self.load(challenge_id).await
302 else {
303 return;
304 };
305 if challenge.status != ChallengeStatus::Processing {
306 return;
307 }
308
309 let orders = OrderService {
310 database: &self.database,
311 audit: &self.audit,
312 profile: &profile,
313 };
314 if orders
315 .abandon_validation(&mut challenge, &mut authz, &mut order, reason)
316 .await
317 .is_err()
318 {
319 error!(
320 event = "challenge_validation_abandon_failed",
321 outcome = "failure",
322 challenge_id = %challenge_id,
323 "the challenge could not be marked invalid and will stay `processing` \
324 until its authorization expires"
325 );
326 }
327 }
328
329 async fn recover(&self, queue: &JobQueue) {
340 let stranded = match Challenge::find_processing(now_secs(), &self.database).await {
341 Ok(rows) => rows,
342 Err(error) => {
343 error!(
344 event = "challenge_validation_recovery_failed",
345 outcome = "failure",
346 error = %error
347 );
348 return;
349 }
350 };
351 for (challenge_id, authz_expires) in stranded {
352 queue
353 .enqueue_or_log(challenge_validate_spec(
354 challenge_id.to_string().as_str(),
355 None,
356 authz_expires,
357 ))
358 .await;
359 }
360 }
361}
362
363const LEASE_HEADROOM: Duration = Duration::from_secs(30);
368
369#[cfg(test)]
370mod tests {
371 use super::*;
372 use crate::acme::order::tests::{account, profile};
373 use acme_proxy_core::identifier::Identifier;
374 use acme_proxy_net::challenge::ChallengeError;
375 use acme_proxy_net::challenge::ChallengeRegistry;
376 use acme_proxy_net::challenge::ChallengeValidator;
377 use acme_proxy_net::challenge::ValidationContext;
378 use acme_proxy_store::status::AuthzStatus;
379 use acme_proxy_store::status::OrderStatus;
380
381 struct Refusing;
384
385 #[async_trait::async_trait]
386 impl ChallengeValidator for Refusing {
387 fn typ(&self) -> &'static str {
388 "http-01"
389 }
390 async fn validate(&self, _ctx: &ValidationContext<'_>) -> Result<(), ChallengeError> {
391 Err(ChallengeError::IncorrectResponse("wrong body".into()))
392 }
393 }
394
395 async fn subject(
397 database: &Arc<Database>,
398 account: &Account,
399 ) -> (Order, Authorization, Challenge) {
400 let order = Order::create(
401 "default",
402 account.id,
403 acme_proxy_store::testutil::dns_identifiers(&["a.example.com"]),
404 now_secs() + 3600,
405 None,
406 None,
407 database,
408 )
409 .await
410 .unwrap();
411 let authz = Authorization::create(
412 order.id,
413 Identifier::dns("a.example.com"),
414 order.expires,
415 database,
416 )
417 .await
418 .unwrap();
419 let challenge = Challenge::create(authz.id, "http-01", database)
420 .await
421 .unwrap();
422 (order, authz, challenge)
423 }
424
425 fn handler(database: &Arc<Database>, challenges: ChallengeRegistry) -> ChallengeValidateJob {
427 let profile = Arc::new(profile(database, challenges));
428 ChallengeValidateJob::new(
429 database.clone(),
430 Arc::new(Auditor::offline(database.clone())),
431 vec![(profile.name.clone(), profile)],
432 )
433 }
434
435 fn row(challenge_id: &str) -> Job {
436 let spec = challenge_validate_spec(challenge_id, None, now_secs() + 3600);
437 Job {
438 dedup_key: spec.key.clone(),
439 payload: spec.payload.clone(),
440 ..acme_proxy_store::testutil::job_fixture()
441 }
442 }
443
444 #[tokio::test]
445 async fn a_claimed_challenge_is_validated_and_its_order_promoted() {
446 let database = Arc::new(Database::connect_in_memory().await.unwrap());
447 let account = account(&database).await;
448 let (order, authz, mut challenge) = subject(&database, &account).await;
449 assert_eq!(
450 challenge.claim_for_validation(0, &database).await.unwrap(),
451 acme_proxy_store::authz::ValidationClaim::Claimed
452 );
453
454 let job = handler(&database, ChallengeRegistry::default());
456 assert!(matches!(
457 job.run(&row(challenge.id.to_string().as_str())).await,
458 JobOutcome::Done
459 ));
460
461 let reloaded = Challenge::find_by_id(challenge.id.to_string().as_str(), &database)
462 .await
463 .unwrap()
464 .unwrap();
465 assert_eq!(reloaded.status, ChallengeStatus::Valid);
466 assert_eq!(
467 Authorization::find_by_id(authz.id.to_string().as_str(), &database)
468 .await
469 .unwrap()
470 .unwrap()
471 .status,
472 AuthzStatus::Valid
473 );
474 assert_eq!(
475 Order::find_by_id(order.id.to_string().as_str(), &database)
476 .await
477 .unwrap()
478 .unwrap()
479 .status,
480 OrderStatus::Ready
481 );
482 }
483
484 #[tokio::test]
488 async fn a_refused_validation_is_recorded_and_the_job_is_done() {
489 let database = Arc::new(Database::connect_in_memory().await.unwrap());
490 let account = account(&database).await;
491 let (order, _, mut challenge) = subject(&database, &account).await;
492 assert_eq!(
493 challenge.claim_for_validation(0, &database).await.unwrap(),
494 acme_proxy_store::authz::ValidationClaim::Claimed
495 );
496
497 let registry = ChallengeRegistry::new(
498 vec![Arc::new(Refusing)],
499 vec!["http-01".to_string()],
500 false,
501 Duration::from_secs(5),
502 );
503 let job = handler(&database, registry);
504 assert!(matches!(
505 job.run(&row(challenge.id.to_string().as_str())).await,
506 JobOutcome::Done
507 ));
508
509 assert_eq!(
510 Order::find_by_id(order.id.to_string().as_str(), &database)
511 .await
512 .unwrap()
513 .unwrap()
514 .status,
515 OrderStatus::Invalid
516 );
517 }
518
519 #[tokio::test]
520 async fn an_unmounted_profile_is_retried() {
521 let database = Arc::new(Database::connect_in_memory().await.unwrap());
522 let account = account(&database).await;
523 let (_, _, mut challenge) = subject(&database, &account).await;
524 assert_eq!(
525 challenge.claim_for_validation(0, &database).await.unwrap(),
526 acme_proxy_store::authz::ValidationClaim::Claimed
527 );
528
529 let job = ChallengeValidateJob::new(
531 database.clone(),
532 Arc::new(Auditor::offline(database.clone())),
533 vec![],
534 );
535 let JobOutcome::Retry(reason) = job.run(&row(challenge.id.to_string().as_str())).await
536 else {
537 panic!("an unmounted profile must be retried, never failed");
538 };
539 assert!(reason.contains("not mounted"), "{reason}");
540 }
541
542 #[tokio::test]
543 async fn a_vanished_challenge_is_failed() {
544 let database = Arc::new(Database::connect_in_memory().await.unwrap());
545 let job = handler(&database, ChallengeRegistry::default());
546 let missing = uuid::Uuid::now_v7().to_string();
547 assert!(matches!(
548 job.run(&row(&missing)).await,
549 JobOutcome::Failed(_)
550 ));
551 }
552
553 #[tokio::test]
558 async fn a_vanished_parent_is_failed_and_an_unreadable_one_retried() {
559 for (table, reason, delete, drop) in [
561 (
562 "authorizations",
563 "authorization",
564 "DELETE FROM authorizations;",
565 "DROP TABLE authorizations;",
566 ),
567 (
568 "orders",
569 "order",
570 "DELETE FROM orders;",
571 "DROP TABLE orders;",
572 ),
573 (
574 "accounts",
575 "account",
576 "DELETE FROM accounts;",
577 "DROP TABLE accounts;",
578 ),
579 ] {
580 for drop_table in [false, true] {
581 let database = Arc::new(Database::connect_in_memory().await.unwrap());
582 let account = account(&database).await;
583 let (_, _, challenge) = subject(&database, &account).await;
584 let target = if drop_table { drop } else { delete };
587 for statement in ["PRAGMA foreign_keys = OFF;", target] {
588 sqlx::query(statement)
589 .execute(database.raw_pool())
590 .await
591 .unwrap();
592 }
593
594 let job = handler(&database, ChallengeRegistry::default());
595 let outcome = job.run(&row(challenge.id.to_string().as_str())).await;
596 match (drop_table, outcome) {
597 (false, JobOutcome::Failed(why)) => {
598 assert!(why.contains(reason), "{table}: {why}");
599 }
600 (true, JobOutcome::Retry(why)) => {
601 assert!(why.contains(reason), "{table}: {why}");
602 }
603 (_, other) => panic!("{table}, dropped {drop_table}: {other:?}"),
604 }
605 }
606 }
607 }
608
609 #[tokio::test]
610 async fn a_payload_naming_no_challenge_is_failed() {
611 let database = Arc::new(Database::connect_in_memory().await.unwrap());
612 let job = handler(&database, ChallengeRegistry::default());
613 let row = Job {
614 payload: serde_json::json!({}),
615 ..acme_proxy_store::testutil::job_fixture()
616 };
617 assert!(matches!(job.run(&row).await, JobOutcome::Failed(_)));
618 }
619
620 #[tokio::test]
623 async fn a_settled_challenge_is_done_without_revalidating() {
624 let database = Arc::new(Database::connect_in_memory().await.unwrap());
625 let account = account(&database).await;
626 let (_, _, mut challenge) = subject(&database, &account).await;
627 assert_eq!(
628 challenge.claim_for_validation(0, &database).await.unwrap(),
629 acme_proxy_store::authz::ValidationClaim::Claimed
630 );
631
632 let job = handler(&database, ChallengeRegistry::default());
633 assert!(matches!(
634 job.run(&row(challenge.id.to_string().as_str())).await,
635 JobOutcome::Done
636 ));
637
638 let registry = ChallengeRegistry::new(
640 vec![Arc::new(Refusing)],
641 vec!["http-01".to_string()],
642 false,
643 Duration::from_secs(5),
644 );
645 let job = handler(&database, registry);
646 assert!(matches!(
647 job.run(&row(challenge.id.to_string().as_str())).await,
648 JobOutcome::Done
649 ));
650 assert_eq!(
651 Challenge::find_by_id(challenge.id.to_string().as_str(), &database)
652 .await
653 .unwrap()
654 .unwrap()
655 .status,
656 ChallengeStatus::Valid
657 );
658 }
659
660 #[tokio::test]
661 async fn abandoning_marks_the_challenge_and_its_order_invalid() {
662 let database = Arc::new(Database::connect_in_memory().await.unwrap());
663 let account = account(&database).await;
664 let (order, authz, mut challenge) = subject(&database, &account).await;
665 assert_eq!(
666 challenge.claim_for_validation(0, &database).await.unwrap(),
667 acme_proxy_store::authz::ValidationClaim::Claimed
668 );
669
670 let job = handler(&database, ChallengeRegistry::default());
671 job.abandon(
672 &row(challenge.id.to_string().as_str()),
673 "the attempts ran out",
674 )
675 .await;
676
677 let reloaded = Challenge::find_by_id(challenge.id.to_string().as_str(), &database)
678 .await
679 .unwrap()
680 .unwrap();
681 assert_eq!(reloaded.status, ChallengeStatus::Invalid);
682 assert_eq!(
683 Authorization::find_by_id(authz.id.to_string().as_str(), &database)
684 .await
685 .unwrap()
686 .unwrap()
687 .status,
688 AuthzStatus::Invalid
689 );
690 assert_eq!(
691 Order::find_by_id(order.id.to_string().as_str(), &database)
692 .await
693 .unwrap()
694 .unwrap()
695 .status,
696 OrderStatus::Invalid
697 );
698 }
699
700 #[tokio::test]
703 async fn recover_requeues_a_stranded_claim_exactly_once() {
704 let database = Arc::new(Database::connect_in_memory().await.unwrap());
705 let account = account(&database).await;
706 let (_, _, mut challenge) = subject(&database, &account).await;
707 assert_eq!(
708 challenge.claim_for_validation(0, &database).await.unwrap(),
709 acme_proxy_store::authz::ValidationClaim::Claimed
710 );
711
712 let queue = acme_proxy_jobs::testutil::idle_job_queue(database.clone());
713 let job = handler(&database, ChallengeRegistry::default());
714
715 job.recover(&queue).await;
716 assert_eq!(
717 acme_proxy_store::job::Job::count_live(CHALLENGE_VALIDATE_KIND, &database)
718 .await
719 .unwrap(),
720 1
721 );
722
723 job.recover(&queue).await;
725 assert_eq!(
726 acme_proxy_store::job::Job::count_live(CHALLENGE_VALIDATE_KIND, &database)
727 .await
728 .unwrap(),
729 1
730 );
731 }
732
733 #[tokio::test]
735 async fn recover_ignores_a_challenge_nobody_claimed() {
736 let database = Arc::new(Database::connect_in_memory().await.unwrap());
737 let account = account(&database).await;
738 let (_, _, _challenge) = subject(&database, &account).await;
739
740 let queue = acme_proxy_jobs::testutil::idle_job_queue(database.clone());
741 handler(&database, ChallengeRegistry::default())
742 .recover(&queue)
743 .await;
744 assert_eq!(
745 acme_proxy_store::job::Job::count_live(CHALLENGE_VALIDATE_KIND, &database)
746 .await
747 .unwrap(),
748 0
749 );
750 }
751}