use super::*;
use crate::jobs::{JobHandler, JobOutcome, JobRegistry};
use crate::signer::relay::flow::{RELAY_JOB_KIND, RelayJob};
use crate::sqlite::job::Job;
use crate::sqlite::status::OrderStatus;
struct TwoUpstreams {
bypassing: RelaySigner,
challenged: RelaySigner,
tokens: Arc<StubTokens>,
bypassing_upstream: Upstream,
challenged_upstream: Upstream,
chain_a: String,
chain_b: String,
_dirs: (TempDir, TempDir),
}
async fn two_upstreams(db: &Arc<Database>, queue: &crate::jobs::JobQueue) -> TwoUpstreams {
let chain_a = real_chain().await;
let chain_b = real_chain().await;
assert_ne!(
chain_a, chain_b,
"the two upstreams must issue distinguishable chains, or nothing below proves anything"
);
let bypassing_upstream = testsrv::start(Script {
chain: chain_a.clone(),
..Script::default()
})
.await;
let challenged_upstream = testsrv::start(Script {
chain: chain_b.clone(),
pose_challenge: true,
offer_http01: true,
..Script::default()
})
.await;
let dir_a = TempDir::new("upstream-a");
let dir_b = TempDir::new("upstream-b");
let parts = relay_parts(db.clone(), no_notifiers(), queue.clone());
let bypassing = RelaySigner::from_config(
&config(&bypassing_upstream, &dir_a),
&parts,
&crate::signer::CarriedState::new(),
)
.unwrap();
let tokens = Arc::new(StubTokens::default());
let challenged = with_tokens(
RelaySigner::from_config(
&config(&challenged_upstream, &dir_b),
&parts,
&crate::signer::CarriedState::new(),
)
.unwrap(),
tokens.clone(),
);
TwoUpstreams {
bypassing,
challenged,
tokens,
bypassing_upstream,
challenged_upstream,
chain_a,
chain_b,
_dirs: (dir_a, dir_b),
}
}
fn handler(db: &Arc<Database>, pair: &TwoUpstreams) -> RelayJob {
RelayJob::new(
db.clone(),
vec![
("a".to_string(), pair.bypassing.relay_state().unwrap()),
("b".to_string(), pair.challenged.relay_state().unwrap()),
],
)
}
#[tokio::test(flavor = "multi_thread")]
async fn two_relay_backends_register_one_handler() {
let db = database().await;
let pair = two_upstreams(&db, &test_queue(db.clone())).await;
let mut registry = JobRegistry::new();
registry
.register(Arc::new(handler(&db, &pair)))
.expect("one handler over both relay backends must register");
assert_eq!(
registry.kinds(),
vec![RELAY_JOB_KIND],
"two backends contribute one kind between them"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn each_profile_is_relayed_by_its_own_backend() {
let db = database().await;
let queue = test_queue(db.clone());
let pair = two_upstreams(&db, &queue).await;
let mut registry = JobRegistry::new();
registry.register(Arc::new(handler(&db, &pair))).unwrap();
let (shutdown, receiver) = tokio::sync::watch::channel(false);
crate::jobs::spawn_runner(
queue,
Arc::new(registry),
&test_jobs_config(),
receiver.clone(),
);
let order_a = ready_order_for("a", db.clone()).await;
let order_b = ready_order_for("b", db.clone()).await;
for order in [&order_a, &order_b] {
let signer = if order.profile == "a" {
&pair.bypassing
} else {
&pair.challenged
};
signer
.issue(
order.id.to_string().as_str(),
&csr_der(),
&identifiers(),
RequestedValidity::default(),
)
.await
.unwrap();
}
let settled_a = await_status(
db.clone(),
order_a.id.to_string().as_str(),
OrderStatus::Valid,
)
.await;
let settled_b = await_status(
db.clone(),
order_b.id.to_string().as_str(),
OrderStatus::Valid,
)
.await;
let _ = shutdown.send(true);
assert_eq!(
settled_a.certificate.as_deref(),
Some(pair.chain_a.as_str()),
"profile `a` must hold what its own upstream issued"
);
assert_eq!(
settled_b.certificate.as_deref(),
Some(pair.chain_b.as_str()),
"profile `b` must hold what its own upstream issued"
);
assert_eq!(
pair.tokens.published().len(),
1,
"profile `b`'s own backend must be the one that answered its challenge"
);
assert_eq!(pair.challenged_upstream.challenge_triggered(), 1);
assert_eq!(pair.bypassing_upstream.challenge_triggered(), 0);
}
fn job(payload: serde_json::Value) -> Job {
Job {
id: crate::sqlite::id::mint(),
kind: RELAY_JOB_KIND.to_string(),
dedup_key: "ord-1".to_string(),
payload,
status: "running".to_string(),
run_at: now_secs(),
attempts: 1,
max_attempts: 5,
deadline: None,
lease_until: None,
lease_owner: None,
last_error: None,
created_at: now_secs(),
updated_at: now_secs(),
}
}
#[tokio::test(flavor = "multi_thread")]
async fn the_lease_is_the_owning_profile_s_own_poll_timeout() {
let db = database().await;
let queue = test_queue(db.clone());
let chain = real_chain().await;
let quick_upstream = testsrv::start(Script {
chain: chain.clone(),
..Script::default()
})
.await;
let patient_upstream = testsrv::start(Script {
chain,
..Script::default()
})
.await;
let (dir_quick, dir_patient) = (TempDir::new("quick"), TempDir::new("patient"));
let parts = relay_parts(db.clone(), no_notifiers(), queue);
let mut quick_config = config(&quick_upstream, &dir_quick);
quick_config.poll_timeout_secs = 11;
let mut patient_config = config(&patient_upstream, &dir_patient);
patient_config.poll_timeout_secs = 97;
let quick =
RelaySigner::from_config(&quick_config, &parts, &crate::signer::CarriedState::new())
.unwrap();
let patient =
RelaySigner::from_config(&patient_config, &parts, &crate::signer::CarriedState::new())
.unwrap();
let handler = RelayJob::new(
db.clone(),
vec![
("quick".to_string(), quick.relay_state().unwrap()),
("patient".to_string(), patient.relay_state().unwrap()),
],
);
assert_eq!(
handler.lease(&job(
serde_json::json!({"order_id": "o", "profile": "quick"})
)),
Some(Duration::from_secs(11))
);
assert_eq!(
handler.lease(&job(
serde_json::json!({"order_id": "o", "profile": "patient"})
)),
Some(Duration::from_secs(97))
);
assert_eq!(
handler.lease(&job(serde_json::json!({"order_id": "o"}))),
Some(Duration::from_secs(97))
);
}
#[tokio::test(flavor = "multi_thread")]
async fn a_row_for_an_unmounted_profile_is_retried_rather_than_abandoned() {
let db = database().await;
let pair = two_upstreams(&db, &test_queue(db.clone())).await;
let handler = handler(&db, &pair);
match handler
.run(&job(
serde_json::json!({"order_id": super::order_id("ord-1"), "profile": "gone"}),
))
.await
{
JobOutcome::Retry(reason) => assert!(
reason.contains("gone"),
"the retry must name the profile: {reason}"
),
other => panic!("expected a retry, got {other:?}"),
}
}
#[tokio::test(flavor = "multi_thread")]
async fn a_payload_from_before_the_profile_was_recorded_resolves_from_the_order() {
let db = database().await;
let pair = two_upstreams(&db, &test_queue(db.clone())).await;
let handler = handler(&db, &pair);
let order = ready_order_for("b", db.clone()).await;
match handler
.run(&job(serde_json::json!({"order_id": order.id})))
.await
{
JobOutcome::Failed(reason) => assert!(
reason.contains("no upstream order"),
"the backend was resolved and then found no mapping: {reason}"
),
other => panic!("expected a permanent failure, got {other:?}"),
}
match handler
.run(&job(
serde_json::json!({"order_id": super::order_id("ord-vanished")}),
))
.await
{
JobOutcome::Failed(reason) => assert!(reason.contains("no longer exists"), "{reason}"),
other => panic!("expected a permanent failure, got {other:?}"),
}
let elsewhere = ready_order_for("c", db.clone()).await;
match handler
.run(&job(serde_json::json!({"order_id": elsewhere.id})))
.await
{
JobOutcome::Retry(reason) => assert!(reason.contains('c'), "{reason}"),
other => panic!("expected a retry, got {other:?}"),
}
db.pool.close().await;
match handler
.run(&job(serde_json::json!({"order_id": order.id})))
.await
{
JobOutcome::Retry(reason) => {
assert!(reason.contains("reading the local order"), "{reason}")
}
other => panic!("expected a retry, got {other:?}"),
}
}
#[tokio::test(flavor = "multi_thread")]
async fn recovery_re_queues_the_in_flight_orders_of_every_backend() {
let db = database().await;
let queue = test_queue(db.clone());
let pair = two_upstreams(&db, &queue).await;
let mut queued = Vec::new();
for profile in ["a", "b"] {
let order = ready_order_for(profile, db.clone()).await;
UpstreamOrder::create(
order.id.to_string().as_str(),
"https://upstream.example/order/1",
Some("https://upstream.example/order/1/finalize"),
&csr_der(),
&db,
)
.await
.unwrap();
queued.push((profile, order.id));
}
handler(&db, &pair).recover(&queue).await;
for (profile, id) in queued {
let row = crate::sqlite::job::Job::find_live(RELAY_JOB_KIND, id.to_string().as_str(), &db)
.await
.unwrap()
.unwrap_or_else(|| panic!("order {id} on profile `{profile}` must be re-queued"));
assert_eq!(
row.payload.get("profile").and_then(|value| value.as_str()),
Some(profile),
"a recovered row names the profile it belongs to, like a fresh one"
);
}
}