rhi

Coordinated trade for connected markets
git clone https://radroots.dev/git/rhi.git
Log | Files | Refs | README | LICENSE

reconciliation_attempt.rs (25026B)


      1 //! Pure bounded per-source reconciliation-attempt planning and result inventory.
      2 
      3 use core::fmt;
      4 use std::error::Error;
      5 
      6 use sha2::{Digest, Sha256};
      7 
      8 use crate::{
      9     RHI_TRADE_SOURCE_RESULT_MAX_BYTES, RHI_TRADE_SOURCE_RESULT_MAX_EVENTS, RhiConfigDocumentV1,
     10     RhiEvidencePolicyDigest, RhiReconciliationJobId, RhiReconciliationJobPolicy,
     11     RhiReconciliationJobState, RhiReconciliationLease, RhiReconciliationUnixMilliseconds,
     12     RhiTradeSourceCompletion, state_metadata,
     13 };
     14 
     15 /// Exact version of the per-source reconciliation-attempt contract.
     16 pub const RHI_RECONCILIATION_ATTEMPT_CONTRACT_VERSION: u32 = 1;
     17 
     18 /// Maximum configured sources represented by one attempt.
     19 pub const RHI_RECONCILIATION_ATTEMPT_MAX_SOURCES: usize = 16;
     20 
     21 const ATTEMPT_ID_DOMAIN: &[u8] = b"radroots.rhi.reconciliation_attempt.v1\0";
     22 const SELECTOR_DIGEST_DOMAIN: &[u8] = b"radroots.rhi.reconciliation_selector.v1\0";
     23 const REQUEST_ID_DOMAIN: &[u8] = b"radroots.rhi.reconciliation_source_request.v1\0";
     24 const SOURCE_SELECTOR: &str = "trade_mutation_lineage_v1";
     25 const SOURCE_KIND: &str = "nostr_relay";
     26 const EVENT_KIND_COUNT: u32 = 5;
     27 const EVENT_KINDS: [u32; 5] = [3470, 3471, 3472, 3473, 3474];
     28 const _: [(); EVENT_KIND_COUNT as usize] = [(); EVENT_KINDS.len()];
     29 const MAX_UNIX_MILLISECONDS: u64 = i64::MAX as u64;
     30 
     31 /// Stable source-free attempt-model failure class.
     32 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
     33 pub enum RhiReconciliationAttemptErrorKind {
     34     InvalidInput,
     35     InvalidConfiguration,
     36     PolicyMismatch,
     37     LeaseExpired,
     38     ResultInventory,
     39 }
     40 
     41 impl RhiReconciliationAttemptErrorKind {
     42     /// Returns the stable machine-readable failure code.
     43     #[must_use]
     44     pub const fn code(self) -> &'static str {
     45         match self {
     46             Self::InvalidInput => "reconciliation_attempt_input_invalid",
     47             Self::InvalidConfiguration => "reconciliation_attempt_configuration_invalid",
     48             Self::PolicyMismatch => "reconciliation_attempt_policy_mismatch",
     49             Self::LeaseExpired => "reconciliation_attempt_lease_expired",
     50             Self::ResultInventory => "reconciliation_attempt_result_inventory_invalid",
     51         }
     52     }
     53 }
     54 
     55 /// Redacted source-free attempt-model failure.
     56 #[derive(Clone, Copy, PartialEq, Eq)]
     57 pub struct RhiReconciliationAttemptError {
     58     kind: RhiReconciliationAttemptErrorKind,
     59 }
     60 
     61 impl RhiReconciliationAttemptError {
     62     const fn new(kind: RhiReconciliationAttemptErrorKind) -> Self {
     63         Self { kind }
     64     }
     65 
     66     /// Returns the stable failure class.
     67     #[must_use]
     68     pub const fn kind(self) -> RhiReconciliationAttemptErrorKind {
     69         self.kind
     70     }
     71 
     72     /// Returns the stable machine-readable failure code.
     73     #[must_use]
     74     pub const fn code(self) -> &'static str {
     75         self.kind.code()
     76     }
     77 }
     78 
     79 impl fmt::Display for RhiReconciliationAttemptError {
     80     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
     81         formatter.write_str(match self.kind {
     82             RhiReconciliationAttemptErrorKind::InvalidInput => {
     83                 "RHI reconciliation attempt input is invalid"
     84             }
     85             RhiReconciliationAttemptErrorKind::InvalidConfiguration => {
     86                 "RHI reconciliation attempt configuration is invalid"
     87             }
     88             RhiReconciliationAttemptErrorKind::PolicyMismatch => {
     89                 "RHI reconciliation attempt policy does not match the claimed job"
     90             }
     91             RhiReconciliationAttemptErrorKind::LeaseExpired => {
     92                 "RHI reconciliation attempt lease has expired"
     93             }
     94             RhiReconciliationAttemptErrorKind::ResultInventory => {
     95                 "RHI reconciliation attempt result inventory is invalid"
     96             }
     97         })
     98     }
     99 }
    100 
    101 impl fmt::Debug for RhiReconciliationAttemptError {
    102     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    103         formatter
    104             .debug_struct("RhiReconciliationAttemptError")
    105             .field("kind", &self.kind)
    106             .finish()
    107     }
    108 }
    109 
    110 impl Error for RhiReconciliationAttemptError {}
    111 
    112 macro_rules! redacted_digest {
    113     ($name:ident, $documentation:literal) => {
    114         #[doc = $documentation]
    115         #[derive(Clone, Copy, PartialEq, Eq, Hash)]
    116         pub struct $name([u8; 32]);
    117 
    118         impl $name {
    119             /// Returns the exact identity bytes.
    120             #[must_use]
    121             pub const fn as_bytes(&self) -> &[u8; 32] {
    122                 &self.0
    123             }
    124         }
    125 
    126         impl fmt::Debug for $name {
    127             fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    128                 formatter.write_str(concat!(stringify!($name), "([redacted])"))
    129             }
    130         }
    131     };
    132 }
    133 
    134 redacted_digest!(
    135     RhiReconciliationAttemptId,
    136     "Domain-separated identity of one claimed reconciliation attempt."
    137 );
    138 redacted_digest!(
    139     RhiReconciliationSourceRequestId,
    140     "Domain-separated identity of one exact source request."
    141 );
    142 redacted_digest!(
    143     RhiReconciliationSourceSelectorDigest,
    144     "Domain-separated digest of the exact trade-mutation selector."
    145 );
    146 
    147 /// One exact configured source request within a claimed job attempt.
    148 #[derive(Clone, PartialEq, Eq)]
    149 pub struct RhiReconciliationSourceRequest {
    150     id: RhiReconciliationSourceRequestId,
    151     source_id: Box<str>,
    152     trade_id: radroots_event::id::TradeId,
    153     required: bool,
    154     selector_digest: RhiReconciliationSourceSelectorDigest,
    155     attempt_started_at: RhiReconciliationUnixMilliseconds,
    156     deadline: RhiReconciliationUnixMilliseconds,
    157     lookback_seconds: u64,
    158     maximum_events: u32,
    159     maximum_bytes: u64,
    160 }
    161 
    162 impl RhiReconciliationSourceRequest {
    163     /// Returns the exact derived request identity.
    164     #[must_use]
    165     pub const fn id(&self) -> RhiReconciliationSourceRequestId {
    166         self.id
    167     }
    168 
    169     /// Returns the validated configured source ID.
    170     #[must_use]
    171     pub fn source_id(&self) -> &str {
    172         &self.source_id
    173     }
    174 
    175     /// Returns the exact trade selected by this request.
    176     #[must_use]
    177     pub const fn trade_id(&self) -> radroots_event::id::TradeId {
    178         self.trade_id
    179     }
    180 
    181     /// Reports whether this source is required by the evidence policy.
    182     #[must_use]
    183     pub const fn required(&self) -> bool {
    184         self.required
    185     }
    186 
    187     /// Returns the exact base-selector digest.
    188     #[must_use]
    189     pub const fn selector_digest(&self) -> RhiReconciliationSourceSelectorDigest {
    190         self.selector_digest
    191     }
    192 
    193     /// Returns the injected attempt start in integer UTC milliseconds.
    194     #[must_use]
    195     pub const fn attempt_started_at(&self) -> RhiReconciliationUnixMilliseconds {
    196         self.attempt_started_at
    197     }
    198 
    199     /// Returns the absolute source deadline capped by attempt and lease expiry.
    200     #[must_use]
    201     pub const fn deadline(&self) -> RhiReconciliationUnixMilliseconds {
    202         self.deadline
    203     }
    204 
    205     /// Returns the configured initial lookback in whole seconds.
    206     #[must_use]
    207     pub const fn lookback_seconds(&self) -> u64 {
    208         self.lookback_seconds
    209     }
    210 
    211     /// Returns the configured maximum accepted event count.
    212     #[must_use]
    213     pub const fn maximum_events(&self) -> u32 {
    214         self.maximum_events
    215     }
    216 
    217     /// Returns the configured maximum accepted original-event bytes.
    218     #[must_use]
    219     pub const fn maximum_bytes(&self) -> u64 {
    220         self.maximum_bytes
    221     }
    222 }
    223 
    224 impl fmt::Debug for RhiReconciliationSourceRequest {
    225     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    226         formatter
    227             .debug_struct("RhiReconciliationSourceRequest")
    228             .field("identity", &"[redacted]")
    229             .field("required", &self.required)
    230             .field("attempt_started_at", &self.attempt_started_at)
    231             .field("deadline", &self.deadline)
    232             .field("lookback_seconds", &self.lookback_seconds)
    233             .field("maximum_events", &self.maximum_events)
    234             .field("maximum_bytes", &self.maximum_bytes)
    235             .finish()
    236     }
    237 }
    238 
    239 /// Canonically ordered per-source plan for one claimed durable job attempt.
    240 #[derive(Clone, PartialEq, Eq)]
    241 pub struct RhiReconciliationAttemptPlan {
    242     id: RhiReconciliationAttemptId,
    243     job_id: RhiReconciliationJobId,
    244     input_generation: u64,
    245     policy_digest: RhiEvidencePolicyDigest,
    246     attempt_started_at: RhiReconciliationUnixMilliseconds,
    247     deadline: RhiReconciliationUnixMilliseconds,
    248     requests: Box<[RhiReconciliationSourceRequest]>,
    249 }
    250 
    251 impl RhiReconciliationAttemptPlan {
    252     /// Derives the complete bounded source plan from one unexpired claimed lease.
    253     pub fn from_claim(
    254         lease: RhiReconciliationLease,
    255         configuration: &RhiConfigDocumentV1,
    256         attempt_started_at: RhiReconciliationUnixMilliseconds,
    257     ) -> Result<Self, RhiReconciliationAttemptError> {
    258         let job = lease.job();
    259         if job.state() != RhiReconciliationJobState::Leased || job.attempt_count() == 0 {
    260             return Err(error(RhiReconciliationAttemptErrorKind::InvalidInput));
    261         }
    262         if attempt_started_at >= lease.lease_expires() {
    263             return Err(error(RhiReconciliationAttemptErrorKind::LeaseExpired));
    264         }
    265         let normalized = configuration.normalized();
    266         let policy_digest = state_metadata::evidence_policy_digest(normalized)
    267             .map_err(|_| error(RhiReconciliationAttemptErrorKind::InvalidConfiguration))?;
    268         if policy_digest != job.evidence_policy_digest() {
    269             return Err(error(RhiReconciliationAttemptErrorKind::PolicyMismatch));
    270         }
    271         let configured_job_policy =
    272             RhiReconciliationJobPolicy::from_configuration(configuration)
    273                 .map_err(|_| error(RhiReconciliationAttemptErrorKind::InvalidConfiguration))?;
    274         if !job.attempt_policy_matches(configured_job_policy) {
    275             return Err(error(RhiReconciliationAttemptErrorKind::PolicyMismatch));
    276         }
    277         let attempt_deadline_ms = bounded_integer(
    278             normalized,
    279             "/reconciliation/attempt_deadline_ms",
    280             100,
    281             30_000,
    282         )?;
    283         let deadline = absolute_deadline(
    284             attempt_started_at,
    285             attempt_deadline_ms,
    286             lease.lease_expires(),
    287         )?;
    288         let maximum_events = bounded_integer(
    289             normalized,
    290             "/resource_limits/source_results/events",
    291             1,
    292             RHI_TRADE_SOURCE_RESULT_MAX_EVENTS as u64,
    293         )
    294         .and_then(|value| {
    295             u32::try_from(value)
    296                 .map_err(|_| error(RhiReconciliationAttemptErrorKind::InvalidConfiguration))
    297         })?;
    298         let maximum_bytes = bounded_integer(
    299             normalized,
    300             "/resource_limits/source_results/bytes",
    301             1,
    302             RHI_TRADE_SOURCE_RESULT_MAX_BYTES as u64,
    303         )?;
    304         let sources = normalized
    305             .pointer("/evidence/sources")
    306             .and_then(serde_json::Value::as_array)
    307             .filter(|sources| {
    308                 !sources.is_empty() && sources.len() <= RHI_RECONCILIATION_ATTEMPT_MAX_SOURCES
    309             })
    310             .ok_or_else(|| error(RhiReconciliationAttemptErrorKind::InvalidConfiguration))?;
    311         let id = attempt_id(job.id(), job.attempt_count());
    312         let selector_digest = selector_digest(job.trade_id());
    313         let mut requests = Vec::with_capacity(sources.len());
    314         let mut prior_source_id: Option<&str> = None;
    315         for source in sources {
    316             let source_id = bounded_source_id(source)?;
    317             if prior_source_id.is_some_and(|prior| prior >= source_id)
    318                 || source.pointer("/kind").and_then(serde_json::Value::as_str) != Some(SOURCE_KIND)
    319                 || source
    320                     .pointer("/selector")
    321                     .and_then(serde_json::Value::as_str)
    322                     != Some(SOURCE_SELECTOR)
    323             {
    324                 return Err(error(
    325                     RhiReconciliationAttemptErrorKind::InvalidConfiguration,
    326                 ));
    327             }
    328             prior_source_id = Some(source_id);
    329             let required = source
    330                 .pointer("/required")
    331                 .and_then(serde_json::Value::as_bool)
    332                 .ok_or_else(|| error(RhiReconciliationAttemptErrorKind::InvalidConfiguration))?;
    333             let source_deadline_ms = bounded_integer(source, "/deadline_ms", 100, 30_000)?;
    334             if source_deadline_ms > attempt_deadline_ms {
    335                 return Err(error(
    336                     RhiReconciliationAttemptErrorKind::InvalidConfiguration,
    337                 ));
    338             }
    339             let source_deadline =
    340                 absolute_deadline(attempt_started_at, source_deadline_ms, deadline)?;
    341             let lookback_seconds = bounded_integer(source, "/lookback_seconds", 60, 2_678_400)?;
    342             requests.push(RhiReconciliationSourceRequest {
    343                 id: request_id(RequestIdentityMaterial {
    344                     attempt_id: id,
    345                     source_id,
    346                     required,
    347                     selector_digest,
    348                     attempt_started_at,
    349                     deadline: source_deadline,
    350                     lookback_seconds,
    351                     maximum_events,
    352                     maximum_bytes,
    353                 }),
    354                 source_id: source_id.into(),
    355                 trade_id: job.trade_id(),
    356                 required,
    357                 selector_digest,
    358                 attempt_started_at,
    359                 deadline: source_deadline,
    360                 lookback_seconds,
    361                 maximum_events,
    362                 maximum_bytes,
    363             });
    364         }
    365         Ok(Self {
    366             id,
    367             job_id: job.id(),
    368             input_generation: job.input_generation(),
    369             policy_digest,
    370             attempt_started_at,
    371             deadline,
    372             requests: requests.into_boxed_slice(),
    373         })
    374     }
    375 
    376     /// Returns the exact derived attempt identity.
    377     #[must_use]
    378     pub const fn id(&self) -> RhiReconciliationAttemptId {
    379         self.id
    380     }
    381 
    382     /// Returns the claimed durable job identity.
    383     #[must_use]
    384     pub const fn job_id(&self) -> RhiReconciliationJobId {
    385         self.job_id
    386     }
    387 
    388     /// Returns the exact dirty generation fenced by the claimed job.
    389     #[must_use]
    390     pub const fn input_generation(&self) -> u64 {
    391         self.input_generation
    392     }
    393 
    394     /// Returns the exact normalized evidence-policy digest.
    395     #[must_use]
    396     pub const fn evidence_policy_digest(&self) -> RhiEvidencePolicyDigest {
    397         self.policy_digest
    398     }
    399 
    400     /// Returns the injected attempt start in integer UTC milliseconds.
    401     #[must_use]
    402     pub const fn attempt_started_at(&self) -> RhiReconciliationUnixMilliseconds {
    403         self.attempt_started_at
    404     }
    405 
    406     /// Returns the absolute attempt deadline capped by lease expiry.
    407     #[must_use]
    408     pub const fn deadline(&self) -> RhiReconciliationUnixMilliseconds {
    409         self.deadline
    410     }
    411 
    412     /// Returns the canonical configured-source request inventory.
    413     #[must_use]
    414     pub fn requests(&self) -> &[RhiReconciliationSourceRequest] {
    415         &self.requests
    416     }
    417 }
    418 
    419 impl fmt::Debug for RhiReconciliationAttemptPlan {
    420     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    421         formatter
    422             .debug_struct("RhiReconciliationAttemptPlan")
    423             .field("identity", &"[redacted]")
    424             .field("input_generation", &self.input_generation)
    425             .field("attempt_started_at", &self.attempt_started_at)
    426             .field("deadline", &self.deadline)
    427             .field("source_count", &self.requests.len())
    428             .finish()
    429     }
    430 }
    431 
    432 /// Bounded source result bound to one exact request identity.
    433 #[derive(Clone, Copy, PartialEq, Eq)]
    434 pub struct RhiReconciliationSourceResult {
    435     request_id: RhiReconciliationSourceRequestId,
    436     outcome: RhiTradeSourceCompletion,
    437     started_at: RhiReconciliationUnixMilliseconds,
    438     finished_at: RhiReconciliationUnixMilliseconds,
    439     accepted_event_count: u32,
    440     accepted_event_bytes: u64,
    441 }
    442 
    443 impl RhiReconciliationSourceResult {
    444     /// Validates exact timing and result bounds for one source request.
    445     pub fn new(
    446         request: &RhiReconciliationSourceRequest,
    447         outcome: RhiTradeSourceCompletion,
    448         started_at: RhiReconciliationUnixMilliseconds,
    449         finished_at: RhiReconciliationUnixMilliseconds,
    450         accepted_event_count: u32,
    451         accepted_event_bytes: u64,
    452     ) -> Result<Self, RhiReconciliationAttemptError> {
    453         let before_deadline = finished_at < request.deadline();
    454         if started_at < request.attempt_started_at()
    455             || started_at > finished_at
    456             || started_at >= request.deadline()
    457             || (outcome == RhiTradeSourceCompletion::IncompleteTimeout && before_deadline)
    458             || (outcome != RhiTradeSourceCompletion::IncompleteTimeout && !before_deadline)
    459             || accepted_event_count > request.maximum_events()
    460             || accepted_event_bytes > request.maximum_bytes()
    461             || (accepted_event_count == 0) != (accepted_event_bytes == 0)
    462             || (outcome == RhiTradeSourceCompletion::Unsupported
    463                 && (accepted_event_count != 0 || accepted_event_bytes != 0))
    464         {
    465             return Err(error(RhiReconciliationAttemptErrorKind::InvalidInput));
    466         }
    467         Ok(Self {
    468             request_id: request.id(),
    469             outcome,
    470             started_at,
    471             finished_at,
    472             accepted_event_count,
    473             accepted_event_bytes,
    474         })
    475     }
    476 
    477     /// Returns the exact request identity this result satisfies.
    478     #[must_use]
    479     pub const fn request_id(self) -> RhiReconciliationSourceRequestId {
    480         self.request_id
    481     }
    482 
    483     /// Returns the stable terminal source-completion classification.
    484     #[must_use]
    485     pub const fn outcome(self) -> RhiTradeSourceCompletion {
    486         self.outcome
    487     }
    488 
    489     /// Returns the injected source-operation start in UTC milliseconds.
    490     #[must_use]
    491     pub const fn started_at(self) -> RhiReconciliationUnixMilliseconds {
    492         self.started_at
    493     }
    494 
    495     /// Returns the injected source-operation finish in UTC milliseconds.
    496     #[must_use]
    497     pub const fn finished_at(self) -> RhiReconciliationUnixMilliseconds {
    498         self.finished_at
    499     }
    500 
    501     /// Returns the bounded accepted event count.
    502     #[must_use]
    503     pub const fn accepted_event_count(self) -> u32 {
    504         self.accepted_event_count
    505     }
    506 
    507     /// Returns the bounded accepted original-event bytes.
    508     #[must_use]
    509     pub const fn accepted_event_bytes(self) -> u64 {
    510         self.accepted_event_bytes
    511     }
    512 }
    513 
    514 impl fmt::Debug for RhiReconciliationSourceResult {
    515     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    516         formatter
    517             .debug_struct("RhiReconciliationSourceResult")
    518             .field("request_id", &"[redacted]")
    519             .field("outcome", &self.outcome)
    520             .field("started_at", &self.started_at)
    521             .field("finished_at", &self.finished_at)
    522             .field("accepted_event_count", &self.accepted_event_count)
    523             .field("accepted_event_bytes", &self.accepted_event_bytes)
    524             .finish()
    525     }
    526 }
    527 
    528 /// Exact canonically ordered result inventory for one attempt plan.
    529 #[derive(Clone, PartialEq, Eq)]
    530 pub struct RhiReconciliationAttemptResults {
    531     attempt_id: RhiReconciliationAttemptId,
    532     results: Box<[RhiReconciliationSourceResult]>,
    533 }
    534 
    535 impl RhiReconciliationAttemptResults {
    536     /// Boundedly ingests exactly one result for each request in canonical order.
    537     pub fn new<I>(
    538         plan: &RhiReconciliationAttemptPlan,
    539         results: I,
    540     ) -> Result<Self, RhiReconciliationAttemptError>
    541     where
    542         I: IntoIterator<Item = RhiReconciliationSourceResult>,
    543     {
    544         let results = results
    545             .into_iter()
    546             .take(plan.requests.len().saturating_add(1))
    547             .collect::<Vec<_>>();
    548         if results.len() != plan.requests.len()
    549             || results
    550                 .iter()
    551                 .zip(plan.requests.iter())
    552                 .any(|(result, request)| result.request_id != request.id)
    553         {
    554             return Err(error(RhiReconciliationAttemptErrorKind::ResultInventory));
    555         }
    556         Ok(Self {
    557             attempt_id: plan.id,
    558             results: results.into_boxed_slice(),
    559         })
    560     }
    561 
    562     /// Returns the exact attempt identity satisfied by this inventory.
    563     #[must_use]
    564     pub const fn attempt_id(&self) -> RhiReconciliationAttemptId {
    565         self.attempt_id
    566     }
    567 
    568     /// Returns the canonical exact result inventory.
    569     #[must_use]
    570     pub fn results(&self) -> &[RhiReconciliationSourceResult] {
    571         &self.results
    572     }
    573 }
    574 
    575 impl fmt::Debug for RhiReconciliationAttemptResults {
    576     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    577         formatter
    578             .debug_struct("RhiReconciliationAttemptResults")
    579             .field("attempt_id", &"[redacted]")
    580             .field("result_count", &self.results.len())
    581             .finish()
    582     }
    583 }
    584 
    585 fn bounded_source_id(source: &serde_json::Value) -> Result<&str, RhiReconciliationAttemptError> {
    586     source
    587         .pointer("/source_id")
    588         .and_then(serde_json::Value::as_str)
    589         .filter(|value| {
    590             !value.is_empty()
    591                 && value.len() <= 64
    592                 && value.bytes().enumerate().all(|(index, byte)| {
    593                     if index == 0 {
    594                         byte.is_ascii_lowercase()
    595                     } else {
    596                         byte.is_ascii_lowercase()
    597                             || byte.is_ascii_digit()
    598                             || matches!(byte, b'_' | b'-')
    599                     }
    600                 })
    601         })
    602         .ok_or_else(|| error(RhiReconciliationAttemptErrorKind::InvalidConfiguration))
    603 }
    604 
    605 fn bounded_integer(
    606     value: &serde_json::Value,
    607     pointer: &str,
    608     minimum: u64,
    609     maximum: u64,
    610 ) -> Result<u64, RhiReconciliationAttemptError> {
    611     value
    612         .pointer(pointer)
    613         .and_then(serde_json::Value::as_u64)
    614         .filter(|value| (minimum..=maximum).contains(value))
    615         .ok_or_else(|| error(RhiReconciliationAttemptErrorKind::InvalidConfiguration))
    616 }
    617 
    618 fn absolute_deadline(
    619     started_at: RhiReconciliationUnixMilliseconds,
    620     duration_ms: u64,
    621     ceiling: RhiReconciliationUnixMilliseconds,
    622 ) -> Result<RhiReconciliationUnixMilliseconds, RhiReconciliationAttemptError> {
    623     let deadline = started_at
    624         .get()
    625         .checked_add(duration_ms)
    626         .filter(|value| *value <= MAX_UNIX_MILLISECONDS)
    627         .map(|value| value.min(ceiling.get()))
    628         .filter(|value| *value > started_at.get())
    629         .ok_or_else(|| error(RhiReconciliationAttemptErrorKind::InvalidInput))?;
    630     RhiReconciliationUnixMilliseconds::new(deadline)
    631         .map_err(|_| error(RhiReconciliationAttemptErrorKind::InvalidInput))
    632 }
    633 
    634 pub(crate) fn attempt_id(
    635     job_id: RhiReconciliationJobId,
    636     attempt_count: u16,
    637 ) -> RhiReconciliationAttemptId {
    638     let mut hasher = Sha256::new();
    639     hasher.update(ATTEMPT_ID_DOMAIN);
    640     hasher.update(job_id.as_bytes());
    641     hasher.update(attempt_count.to_be_bytes());
    642     RhiReconciliationAttemptId(hasher.finalize().into())
    643 }
    644 
    645 fn selector_digest(trade_id: radroots_event::id::TradeId) -> RhiReconciliationSourceSelectorDigest {
    646     let mut hasher = Sha256::new();
    647     hasher.update(SELECTOR_DIGEST_DOMAIN);
    648     hash_framed(&mut hasher, SOURCE_SELECTOR.as_bytes());
    649     hasher.update(EVENT_KIND_COUNT.to_be_bytes());
    650     for kind in EVENT_KINDS {
    651         hasher.update(kind.to_be_bytes());
    652     }
    653     hasher.update(b"d");
    654     hasher.update(trade_id.as_bytes());
    655     RhiReconciliationSourceSelectorDigest(hasher.finalize().into())
    656 }
    657 
    658 struct RequestIdentityMaterial<'source> {
    659     attempt_id: RhiReconciliationAttemptId,
    660     source_id: &'source str,
    661     required: bool,
    662     selector_digest: RhiReconciliationSourceSelectorDigest,
    663     attempt_started_at: RhiReconciliationUnixMilliseconds,
    664     deadline: RhiReconciliationUnixMilliseconds,
    665     lookback_seconds: u64,
    666     maximum_events: u32,
    667     maximum_bytes: u64,
    668 }
    669 
    670 fn request_id(material: RequestIdentityMaterial<'_>) -> RhiReconciliationSourceRequestId {
    671     let mut hasher = Sha256::new();
    672     hasher.update(REQUEST_ID_DOMAIN);
    673     hasher.update(material.attempt_id.as_bytes());
    674     hash_framed(&mut hasher, material.source_id.as_bytes());
    675     hasher.update([u8::from(material.required)]);
    676     hasher.update(material.selector_digest.as_bytes());
    677     hasher.update(material.attempt_started_at.get().to_be_bytes());
    678     hasher.update(material.deadline.get().to_be_bytes());
    679     hasher.update(material.lookback_seconds.to_be_bytes());
    680     hasher.update(material.maximum_events.to_be_bytes());
    681     hasher.update(material.maximum_bytes.to_be_bytes());
    682     RhiReconciliationSourceRequestId(hasher.finalize().into())
    683 }
    684 
    685 fn hash_framed(hasher: &mut Sha256, bytes: &[u8]) {
    686     hasher.update(u64::try_from(bytes.len()).unwrap_or(u64::MAX).to_be_bytes());
    687     hasher.update(bytes);
    688 }
    689 
    690 const fn error(kind: RhiReconciliationAttemptErrorKind) -> RhiReconciliationAttemptError {
    691     RhiReconciliationAttemptError::new(kind)
    692 }