use std::collections::BTreeMap;
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::db::Database;
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, UpstreamChallengeView, UpstreamOrderView};
use super::{ChallengeStrategy, Inner, RelayState, 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 {
targets: BTreeMap<String, Arc<Inner>>,
upstreams: Vec<(Arc<Inner>, Vec<String>)>,
longest_lease: Option<Duration>,
database: Arc<Database>,
}
impl RelayJob {
#[must_use]
pub fn new(database: Arc<Database>, targets: Vec<(String, RelayState)>) -> Self {
let mut upstreams: Vec<(Arc<Inner>, Vec<String>)> = Vec::new();
for (profile, state) in &targets {
match upstreams
.iter_mut()
.find(|(inner, _)| Arc::ptr_eq(inner, &state.0))
{
Some((_, served)) => served.push(profile.clone()),
None => upstreams.push((state.0.clone(), vec![profile.clone()])),
}
}
let longest_lease = targets.iter().map(|(_, state)| state.0.poll.timeout).max();
Self {
targets: targets
.into_iter()
.map(|(profile, state)| (profile, state.0))
.collect(),
upstreams,
longest_lease,
database,
}
}
fn target_for(&self, job: &Job) -> Option<&Arc<Inner>> {
let profile = job.payload.get("profile").and_then(Value::as_str)?;
self.targets.get(profile)
}
async fn resolve(&self, job: &Job, order_id: &str) -> Resolved<'_> {
if let Some(profile) = job.payload.get("profile").and_then(Value::as_str) {
return match self.targets.get(profile) {
Some(inner) => Resolved::Backend(inner),
None => Resolved::Unmounted(profile.to_string()),
};
}
match Order::find_by_id(order_id, &self.database).await {
Ok(Some(order)) => match self.targets.get(&order.profile) {
Some(inner) => Resolved::Backend(inner),
None => Resolved::Unmounted(order.profile),
},
Ok(None) => Resolved::Gone,
Err(error) => Resolved::Unreadable(error.to_string()),
}
}
}
enum Resolved<'a> {
Backend(&'a Arc<Inner>),
Unmounted(String),
Gone,
Unreadable(String),
}
#[async_trait]
impl JobHandler for RelayJob {
fn kind(&self) -> &'static str {
RELAY_JOB_KIND
}
fn lease(&self, job: &Job) -> Option<Duration> {
self.target_for(job)
.map(|inner| inner.poll.timeout)
.or(self.longest_lease)
}
async fn run(&self, job: &Job) -> JobOutcome {
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 inner = match self.resolve(job, order_id).await {
Resolved::Backend(inner) => inner,
Resolved::Unmounted(profile) => {
return JobOutcome::Retry(format!(
"no relay backend is mounted for profile `{profile}`"
));
}
Resolved::Gone => {
return JobOutcome::Failed(format!("local order {order_id} no longer exists"));
}
Resolved::Unreadable(error) => {
return JobOutcome::Retry(format!("reading the local order failed: {error}"));
}
};
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 Some(order_id) = job.payload.get("order_id").and_then(Value::as_str) else {
return;
};
let Resolved::Backend(inner) = self.resolve(job, order_id).await else {
warn!(event = "upstream_relay_backend_unresolved", outcome = "failure", order_id = %order_id, reason = %reason);
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 (actor, client) = relay_actor_and_client(&order, &inner.database).await;
if let Err(error) = abandon_relayed_order(
&mut order,
reason,
actor,
client,
&inner.database,
Some(&inner.metrics),
)
.await
{
error!(event = "upstream_relay_mark_invalid_failed", outcome = "failure", order_id = %order_id, error = %error);
}
}
async fn recover(&self, queue: &JobQueue) {
for (inner, profiles) in &self.upstreams {
let pending = match UpstreamOrder::list_processing(profiles, &inner.database).await {
Ok(pending) => pending,
Err(error) => {
error!(event = "upstream_resume_lookup_failed", outcome = "failure", error = %error);
continue;
}
};
if pending.is_empty() {
continue;
}
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 context = OrderContext::read(row.order_id.to_string().as_str(), inner).await;
queue
.enqueue_or_log(relay_spec(row.order_id.to_string().as_str(), &context))
.await;
}
}
}
}
pub(super) fn relay_spec(order_id: &str, context: &OrderContext) -> JobSpec {
let mut payload = json!({ "order_id": order_id });
if let Some(profile) = &context.profile {
payload["profile"] = json!(profile);
}
JobSpec::now(RELAY_JOB_KIND, order_id)
.with_payload(payload)
.with_deadline(context.deadline)
}
#[derive(Default)]
pub(super) struct OrderContext {
pub(super) deadline: Option<i64>,
pub(super) profile: Option<String>,
}
impl OrderContext {
pub(super) async fn read(order_id: &str, inner: &Inner) -> Self {
match Order::find_by_id(order_id, &inner.database).await {
Ok(Some(order)) => Self {
deadline: Some(order.expires),
profile: Some(order.profile),
},
Ok(None) => Self::default(),
Err(error) => {
warn!(event = "upstream_relay_order_context_lookup_failed", outcome = "failure", order_id = %order_id, error = %error);
Self::default()
}
}
}
}
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 token = require_token(challenge, &authz.identifier.value)?;
let name = crate::challenge::dns_01::record_name(&authz.identifier.value);
let key_authorization = format!("{token}.{thumbprint}");
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 token = require_token(challenge, &authz.identifier.value)?;
let key_authorization = format!("{token}.{thumbprint}");
let _published = http01::PublishedToken::publish(tokens.clone(), 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
.iter()
.find(|challenge| challenge.token.is_some())
.or_else(|| 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(())
}
fn require_token<'a>(
challenge: &'a UpstreamChallengeView,
identifier: &str,
) -> Result<&'a str, RelayFailure> {
challenge.token.as_deref().ok_or_else(|| {
RelayFailure::Permanent(format!(
"upstream {} challenge for {identifier} carries no token",
challenge.typ
))
})
}
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}"));
}
};
let cert_not_after = crate::cert::cert_validity(&leaf)
.ok()
.map(|(_, not_after)| not_after);
if let Err(error) = order
.finalize(
chain,
serial.clone(),
pubkey,
cert_not_after,
&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().to_string(),
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 (actor, client) = relay_actor_and_client(order, &inner.database).await;
crate::audit::AuditRecord::new(event, &order.profile, actor)
.with_order(order)
.with_client(client)
}
async fn relay_actor_and_client(
order: &Order,
database: &Database,
) -> (crate::audit::Actor, crate::audit::ClientContext) {
let mapping = UpstreamOrder::find_by_order_id(&order.id.to_string(), database)
.await
.unwrap_or_else(|error| {
warn!(event = "upstream_order_client_context_lookup_failed", outcome = "failure", order_id = %order.id, error = %error);
None
});
match &mapping {
Some(mapping) => (
crate::audit::Actor::acme(order.account_id.to_string()),
mapping.client(),
),
None => (
crate::audit::Actor::system(),
crate::audit::ClientContext::default(),
),
}
}
pub(crate) async fn abandon_relayed_order(
order: &mut Order,
reason: &str,
actor: crate::audit::Actor,
client: crate::audit::ClientContext,
database: &Database,
metrics: Option<&Arc<crate::metrics::Metrics>>,
) -> Result<(), sqlx::Error> {
let record = crate::audit::AuditRecord::new(
crate::audit::AuditEvent::CertificateIssueFailed,
&order.profile,
actor,
)
.with_order(order)
.with_client(client)
.with_reason("serverInternal")
.with_detail(reason);
if let Some(metrics) = metrics {
metrics.record_audit(&record);
}
crate::audit::write(record, database).await;
let problem = Problem::server_internal("Upstream certificate issuance failed");
order.mark_invalid(problem.to_value(), database).await?;
UpstreamOrder::mark_invalid(&order.id.to_string(), reason, database).await?;
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
#[tokio::test]
async fn abandon_relayed_order_marks_both_rows_and_writes_one_row() {
use crate::sqlite::account::Account;
use crate::sqlite::audit::{AuditEntry, AuditQuery};
use crate::sqlite::db::Database;
use crate::sqlite::order::Identifier;
let database = Database::connect_in_memory().await.unwrap();
let (account, _) = Account::find_or_create(
"default",
&crate::random::random_bytes::<16>(),
Vec::new(),
&crate::audit::ClientContext::default(),
&database,
)
.await
.unwrap();
let mut order = Order::create(
"default",
account.id,
vec![Identifier::dns("a.example.com")],
crate::sqlite::nonce::now_secs() + 3600,
None,
None,
&database,
)
.await
.unwrap();
UpstreamOrder::create(
order.id.to_string().as_str(),
"https://up.example/o/1",
None,
b"csr",
&database,
)
.await
.unwrap();
abandon_relayed_order(
&mut order,
"cancelled by operator",
crate::audit::Actor::admin("root"),
crate::audit::ClientContext {
ip: Some("203.0.113.9".to_string()),
..crate::audit::ClientContext::default()
},
&database,
None,
)
.await
.unwrap();
let reloaded = Order::find_by_id(order.id.to_string().as_str(), &database)
.await
.unwrap()
.unwrap();
assert_eq!(reloaded.status.as_str(), "invalid");
assert!(
reloaded
.error
.as_ref()
.unwrap()
.to_string()
.contains("Upstream certificate issuance failed")
);
assert!(
!reloaded
.error
.as_ref()
.unwrap()
.to_string()
.contains("cancelled by operator")
);
let mapping = UpstreamOrder::find_by_order_id(order.id.to_string().as_str(), &database)
.await
.unwrap()
.unwrap();
assert_eq!(mapping.status, "invalid");
assert_eq!(mapping.error.as_deref(), Some("cancelled by operator"));
let (rows, _) = AuditEntry::search(
&AuditQuery {
limit: 50,
..AuditQuery::default()
},
&database,
)
.await
.unwrap();
assert_eq!(rows.len(), 1);
assert_eq!(rows[0].event, "certificate_issue_failed");
assert_eq!(rows[0].actor_kind, "admin");
assert_eq!(rows[0].actor_id.as_deref(), Some("root"));
assert_eq!(rows[0].client_ip.as_deref(), Some("203.0.113.9"));
}
#[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",
&OrderContext {
deadline: Some(1_234),
profile: Some("le".to_string()),
},
);
assert_eq!(spec.kind, RELAY_JOB_KIND);
assert_eq!(spec.key, "ord-1");
assert_eq!(spec.payload, json!({"order_id": "ord-1", "profile": "le"}));
assert_eq!(spec.deadline, Some(1_234));
}
#[test]
fn a_spec_for_an_unreadable_order_names_neither_profile_nor_deadline() {
let spec = relay_spec("ord-1", &OrderContext::default());
assert_eq!(spec.payload, json!({"order_id": "ord-1"}));
assert_eq!(spec.deadline, None);
}
}