rhi

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

reconciliation_manifest.rs (18103B)


      1 //! Immutable manifest materialization from one confirmed reconciliation commit.
      2 
      3 use core::{fmt, num::NonZeroU64};
      4 use std::{collections::BTreeMap, error::Error, sync::Arc};
      5 
      6 use radroots_event::id::{EventId, MutationId, TradeId};
      7 use radroots_service_host::UnixTimeSeconds;
      8 use radroots_trade::evidence::{
      9     RadrootsTradeEvidenceManifestObservationV1, RadrootsTradeEvidenceManifestSourceResultV1,
     10     RadrootsTradeEvidenceManifestV1, RadrootsTradeEvidencePolicyDigestV1,
     11     RadrootsTradeEvidenceProvenanceDigestV1, RadrootsTradeEvidenceScopePrerequisitesV1,
     12     RadrootsTradeEvidenceSourceCompletionV1, RadrootsTradeEvidenceSourceIdV1,
     13     RadrootsTradeEvidenceSourceRequirementV1, RadrootsTradeEvidenceSourceResultDigestV1,
     14     RadrootsTradeEvidenceSourceResultV1, RadrootsTradeSignedEventDigestV1,
     15 };
     16 use sha2::{Digest, Sha256};
     17 
     18 use crate::{
     19     RhiReconciliationAttemptId, RhiReconciliationAttemptPlan, RhiReconciliationJobId,
     20     RhiReconciliationSourceCommitOutcome, RhiTradeSourceCompletion,
     21     reconciliation_commit::committed_inventory_digest,
     22     reconciliation_replay::{
     23         RhiReconciliationReplayCommitFact, RhiReconciliationReplayCommitParts,
     24     },
     25 };
     26 
     27 /// Exact version of the RHI reconciliation-manifest materialization contract.
     28 pub const RHI_RECONCILIATION_MANIFEST_CONTRACT_VERSION: u32 = 1;
     29 
     30 const SOURCE_RESULT_DIGEST_DOMAIN: &[u8] =
     31     b"radroots.rhi.reconciliation_manifest_source_result.v1\0";
     32 const PROVENANCE_DIGEST_DOMAIN: &[u8] = b"radroots.rhi.evidence_provenance.v1\0";
     33 const SOURCE_SELECTOR: &[u8] = b"trade_mutation_lineage_v1";
     34 pub(crate) const RHI_REDUCER_MAXIMUM_MUTATIONS: usize = 65_536;
     35 pub(crate) const RHI_REDUCER_MAXIMUM_MUTATION_MATERIAL_BYTES: usize = 134_217_728;
     36 
     37 /// Exact non-source prerequisite state bound into one manifest.
     38 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
     39 pub enum RhiReconciliationScopePrerequisites {
     40     Satisfied,
     41     Unsatisfied,
     42 }
     43 
     44 /// Stable source-free manifest materialization failure class.
     45 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
     46 pub enum RhiReconciliationManifestErrorKind {
     47     InvalidObservationTime,
     48     InvalidCommittedInventory,
     49 }
     50 
     51 impl RhiReconciliationManifestErrorKind {
     52     /// Returns the stable machine-readable failure code.
     53     #[must_use]
     54     pub const fn code(self) -> &'static str {
     55         match self {
     56             Self::InvalidObservationTime => "reconciliation_manifest_observation_time_invalid",
     57             Self::InvalidCommittedInventory => "reconciliation_manifest_inventory_invalid",
     58         }
     59     }
     60 }
     61 
     62 /// Redacted source-free reconciliation-manifest failure.
     63 #[derive(Clone, Copy, PartialEq, Eq)]
     64 pub struct RhiReconciliationManifestError {
     65     kind: RhiReconciliationManifestErrorKind,
     66 }
     67 
     68 impl RhiReconciliationManifestError {
     69     /// Returns the stable failure class.
     70     #[must_use]
     71     pub const fn kind(self) -> RhiReconciliationManifestErrorKind {
     72         self.kind
     73     }
     74 
     75     /// Returns the stable machine-readable failure code.
     76     #[must_use]
     77     pub const fn code(self) -> &'static str {
     78         self.kind.code()
     79     }
     80 }
     81 
     82 impl fmt::Display for RhiReconciliationManifestError {
     83     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
     84         formatter.write_str(match self.kind {
     85             RhiReconciliationManifestErrorKind::InvalidObservationTime => {
     86                 "RHI reconciliation manifest observation time is invalid"
     87             }
     88             RhiReconciliationManifestErrorKind::InvalidCommittedInventory => {
     89                 "RHI committed reconciliation inventory is invalid"
     90             }
     91         })
     92     }
     93 }
     94 
     95 impl fmt::Debug for RhiReconciliationManifestError {
     96     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
     97         formatter
     98             .debug_struct("RhiReconciliationManifestError")
     99             .field("kind", &self.kind)
    100             .finish()
    101     }
    102 }
    103 
    104 impl Error for RhiReconciliationManifestError {}
    105 
    106 /// Sealed immutable evidence manifest derived from one confirmed Step190 commit.
    107 ///
    108 /// Callers cannot forge a manifest by constructing its representation.
    109 ///
    110 /// ```compile_fail
    111 /// use rhi::RhiReconciliationManifest;
    112 ///
    113 /// let _forged = RhiReconciliationManifest { inner: todo!() };
    114 /// ```
    115 pub struct RhiReconciliationManifest {
    116     inner: RadrootsTradeEvidenceManifestV1,
    117     attempt_id: RhiReconciliationAttemptId,
    118     job_id: RhiReconciliationJobId,
    119     reducer_mutations: Box<[RhiReducerMutationMaterial]>,
    120 }
    121 
    122 impl RhiReconciliationManifest {
    123     /// Returns the exact RHI manifest-materialization contract version.
    124     #[must_use]
    125     pub const fn contract_version(&self) -> u32 {
    126         RHI_RECONCILIATION_MANIFEST_CONTRACT_VERSION
    127     }
    128 
    129     /// Returns the exact shared manifest encoding contract ID.
    130     #[must_use]
    131     pub const fn shared_manifest_contract_id(&self) -> &'static str {
    132         self.inner.contract_id()
    133     }
    134 
    135     /// Returns the exact shared manifest encoding contract version.
    136     #[must_use]
    137     pub const fn shared_manifest_contract_version(&self) -> u16 {
    138         self.inner.contract_version()
    139     }
    140 
    141     /// Returns the exact trade selected by the committed attempt.
    142     #[must_use]
    143     pub const fn trade_id(&self) -> &TradeId {
    144         self.inner.trade_id()
    145     }
    146 
    147     /// Returns the exact nonzero dirty generation frozen by the manifest.
    148     #[must_use]
    149     pub const fn trade_generation(&self) -> u64 {
    150         self.inner.trade_generation().get()
    151     }
    152 
    153     /// Returns the explicit observation time in UTC seconds.
    154     #[must_use]
    155     pub const fn observed_at_unix_seconds(&self) -> u64 {
    156         self.inner.observed_at_unix_s()
    157     }
    158 
    159     /// Returns the exact canonical manifest bytes.
    160     #[must_use]
    161     pub fn canonical_bytes(&self) -> &[u8] {
    162         self.inner.canonical_bytes()
    163     }
    164 
    165     /// Returns the domain-separated shared manifest digest bytes.
    166     #[must_use]
    167     pub fn digest(&self) -> [u8; 32] {
    168         *self.inner.digest().as_bytes()
    169     }
    170 
    171     /// Returns the exact configured source count.
    172     #[must_use]
    173     pub fn source_count(&self) -> usize {
    174         self.inner.sources().len()
    175     }
    176 
    177     /// Returns the exact accepted source-observation count.
    178     #[must_use]
    179     pub fn observation_count(&self) -> usize {
    180         self.inner.observations().len()
    181     }
    182 
    183     pub(crate) const fn inner(&self) -> &RadrootsTradeEvidenceManifestV1 {
    184         &self.inner
    185     }
    186 
    187     pub(crate) const fn attempt_id(&self) -> RhiReconciliationAttemptId {
    188         self.attempt_id
    189     }
    190 
    191     pub(crate) const fn job_id(&self) -> RhiReconciliationJobId {
    192         self.job_id
    193     }
    194 
    195     pub(crate) fn reducer_mutations(&self) -> &[RhiReducerMutationMaterial] {
    196         &self.reducer_mutations
    197     }
    198 }
    199 
    200 impl fmt::Debug for RhiReconciliationManifest {
    201     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    202         formatter
    203             .debug_struct("RhiReconciliationManifest")
    204             .field("source_count", &self.source_count())
    205             .field("observation_count", &self.observation_count())
    206             .finish_non_exhaustive()
    207     }
    208 }
    209 
    210 impl RhiReconciliationSourceCommitOutcome {
    211     /// Freezes the exact committed inventory into the governed shared manifest.
    212     ///
    213     /// This consumes the sealed commit outcome so uncommitted replay material
    214     /// cannot be relabelled as durable evidence. The observation time must be
    215     /// at or after every committed source-result completion.
    216     pub fn into_evidence_manifest(
    217         self,
    218         observed_at: UnixTimeSeconds,
    219         prerequisites: RhiReconciliationScopePrerequisites,
    220     ) -> Result<RhiReconciliationManifest, RhiReconciliationManifestError> {
    221         freeze_manifest(self.manifest_material, observed_at, prerequisites)
    222     }
    223 }
    224 
    225 pub(crate) struct RhiCommittedManifestMaterial {
    226     attempt_id: RhiReconciliationAttemptId,
    227     job_id: RhiReconciliationJobId,
    228     trade_id: TradeId,
    229     generation: NonZeroU64,
    230     policy_digest: RadrootsTradeEvidencePolicyDigestV1,
    231     latest_finished_unix_ms: u64,
    232     sources: Box<[RadrootsTradeEvidenceManifestSourceResultV1]>,
    233     observations: Box<[RadrootsTradeEvidenceManifestObservationV1]>,
    234     reducer_mutations: Box<[RhiReducerMutationMaterial]>,
    235 }
    236 
    237 pub(crate) struct RhiReducerMutationMaterial {
    238     pub(crate) mutation_id: MutationId,
    239     pub(crate) event_id: EventId,
    240     pub(crate) canonical_content: Arc<[u8]>,
    241 }
    242 
    243 pub(crate) fn committed_manifest_material(
    244     plan: &RhiReconciliationAttemptPlan,
    245     parts: &[RhiReconciliationReplayCommitParts],
    246 ) -> Result<RhiCommittedManifestMaterial, ()> {
    247     let trade_id = parts.first().ok_or(())?.trade_id;
    248     let generation = NonZeroU64::new(plan.input_generation()).ok_or(())?;
    249     let policy_digest =
    250         RadrootsTradeEvidencePolicyDigestV1::from_bytes(*plan.evidence_policy_digest().as_bytes());
    251     let mut latest_finished_unix_ms = 0_u64;
    252     let mut sources = Vec::with_capacity(parts.len());
    253     let observation_capacity = parts.iter().try_fold(0_usize, |total, part| {
    254         total.checked_add(part.facts.len()).ok_or(())
    255     })?;
    256     let mut observations = Vec::with_capacity(observation_capacity);
    257     let mut reducer_mutations = BTreeMap::<[u8; 32], RhiReducerMutationMaterial>::new();
    258 
    259     for (ordinal, part) in parts.iter().enumerate() {
    260         if part.trade_id != trade_id
    261             || part.policy_digest != *plan.evidence_policy_digest().as_bytes()
    262         {
    263             return Err(());
    264         }
    265         latest_finished_unix_ms = latest_finished_unix_ms.max(part.result.finished_at().get());
    266         let source_id =
    267             RadrootsTradeEvidenceSourceIdV1::parse(part.source_id.as_ref()).map_err(|_| ())?;
    268         let result = RadrootsTradeEvidenceSourceResultV1::new(
    269             if part.required {
    270                 RadrootsTradeEvidenceSourceRequirementV1::Required
    271             } else {
    272                 RadrootsTradeEvidenceSourceRequirementV1::Optional
    273             },
    274             map_completion(part.result.outcome()),
    275             part.result.accepted_event_count(),
    276         )
    277         .map_err(|_| ())?;
    278         let inventory_digest = committed_inventory_digest(part).ok_or(())?;
    279         let result_digest = source_result_digest(plan, ordinal, part, inventory_digest)?;
    280         sources.push(RadrootsTradeEvidenceManifestSourceResultV1::new(
    281             source_id.clone(),
    282             result,
    283             RadrootsTradeEvidenceSourceResultDigestV1::from_bytes(result_digest),
    284         ));
    285         for fact in &part.facts {
    286             observations.push(manifest_observation(source_id.clone(), part, fact)?);
    287             reducer_mutations
    288                 .entry(fact.record.mutation_id)
    289                 .and_modify(|current| {
    290                     if fact.record.event_id < *current.event_id.as_bytes() {
    291                         current.event_id = EventId::from_bytes(fact.record.event_id);
    292                         current.canonical_content = fact.record.canonical_content.clone();
    293                     }
    294                 })
    295                 .or_insert_with(|| RhiReducerMutationMaterial {
    296                     mutation_id: MutationId::from_bytes(fact.record.mutation_id),
    297                     event_id: EventId::from_bytes(fact.record.event_id),
    298                     canonical_content: fact.record.canonical_content.clone(),
    299                 });
    300         }
    301     }
    302     if reducer_mutations.len() > RHI_REDUCER_MAXIMUM_MUTATIONS
    303         || reducer_mutations
    304             .values()
    305             .try_fold(0_usize, |total, material| {
    306                 total.checked_add(material.canonical_content.len())
    307             })
    308             .is_none_or(|total| total > RHI_REDUCER_MAXIMUM_MUTATION_MATERIAL_BYTES)
    309     {
    310         return Err(());
    311     }
    312 
    313     Ok(RhiCommittedManifestMaterial {
    314         attempt_id: plan.id(),
    315         job_id: plan.job_id(),
    316         trade_id,
    317         generation,
    318         policy_digest,
    319         latest_finished_unix_ms,
    320         sources: sources.into_boxed_slice(),
    321         observations: observations.into_boxed_slice(),
    322         reducer_mutations: reducer_mutations
    323             .into_values()
    324             .collect::<Vec<_>>()
    325             .into_boxed_slice(),
    326     })
    327 }
    328 
    329 fn freeze_manifest(
    330     material: RhiCommittedManifestMaterial,
    331     observed_at: UnixTimeSeconds,
    332     prerequisites: RhiReconciliationScopePrerequisites,
    333 ) -> Result<RhiReconciliationManifest, RhiReconciliationManifestError> {
    334     let earliest_observation = material.latest_finished_unix_ms.div_ceil(1_000);
    335     if observed_at.get() < earliest_observation || i64::try_from(observed_at.get()).is_err() {
    336         return Err(error(
    337             RhiReconciliationManifestErrorKind::InvalidObservationTime,
    338         ));
    339     }
    340     let inner = RadrootsTradeEvidenceManifestV1::new(
    341         material.trade_id,
    342         material.generation,
    343         material.policy_digest,
    344         observed_at.get(),
    345         match prerequisites {
    346             RhiReconciliationScopePrerequisites::Satisfied => {
    347                 RadrootsTradeEvidenceScopePrerequisitesV1::Satisfied
    348             }
    349             RhiReconciliationScopePrerequisites::Unsatisfied => {
    350                 RadrootsTradeEvidenceScopePrerequisitesV1::Unsatisfied
    351             }
    352         },
    353         material.sources.into_vec(),
    354         material.observations.into_vec(),
    355     )
    356     .map_err(|_| error(RhiReconciliationManifestErrorKind::InvalidCommittedInventory))?;
    357     Ok(RhiReconciliationManifest {
    358         inner,
    359         attempt_id: material.attempt_id,
    360         job_id: material.job_id,
    361         reducer_mutations: material.reducer_mutations,
    362     })
    363 }
    364 
    365 fn map_completion(value: RhiTradeSourceCompletion) -> RadrootsTradeEvidenceSourceCompletionV1 {
    366     match value {
    367         RhiTradeSourceCompletion::Complete => RadrootsTradeEvidenceSourceCompletionV1::Complete,
    368         RhiTradeSourceCompletion::Unsupported => {
    369             RadrootsTradeEvidenceSourceCompletionV1::Unsupported
    370         }
    371         RhiTradeSourceCompletion::IncompleteTimeout
    372         | RhiTradeSourceCompletion::IncompleteUnavailable
    373         | RhiTradeSourceCompletion::IncompleteResourceLimit
    374         | RhiTradeSourceCompletion::IncompleteUnknown => {
    375             RadrootsTradeEvidenceSourceCompletionV1::Incomplete
    376         }
    377     }
    378 }
    379 
    380 fn source_result_digest(
    381     plan: &RhiReconciliationAttemptPlan,
    382     ordinal: usize,
    383     part: &RhiReconciliationReplayCommitParts,
    384     inventory_digest: [u8; 32],
    385 ) -> Result<[u8; 32], ()> {
    386     let mut digest = Sha256::new();
    387     digest.update(SOURCE_RESULT_DIGEST_DOMAIN);
    388     digest.update(plan.id().as_bytes());
    389     digest.update(u32::try_from(ordinal).map_err(|_| ())?.to_be_bytes());
    390     digest.update(part.request_id.as_bytes());
    391     update_framed(&mut digest, part.source_id.as_bytes())?;
    392     digest.update(part.trade_id.as_bytes());
    393     digest.update([u8::from(part.required)]);
    394     digest.update(part.policy_digest);
    395     digest.update(part.selector_digest);
    396     digest.update(part.replay_id.as_bytes());
    397     update_framed(&mut digest, part.result.outcome().code().as_bytes())?;
    398     digest.update(part.result.started_at().get().to_be_bytes());
    399     digest.update(part.result.finished_at().get().to_be_bytes());
    400     digest.update(part.result.accepted_event_count().to_be_bytes());
    401     digest.update(part.result.accepted_event_bytes().to_be_bytes());
    402     digest.update(inventory_digest);
    403     digest.update(part.duplicate_observations.to_be_bytes());
    404     update_optional_u64(&mut digest, part.first_observed_at.map(|value| value.get()));
    405     update_optional_cursor(&mut digest, part.prior_cursor);
    406     digest.update(part.overlap_seconds.to_be_bytes());
    407     digest.update(part.since_unix_seconds.to_be_bytes());
    408     update_optional_cursor(&mut digest, part.cursor_candidate);
    409     digest.update([u8::from(part.eligible_cursor.is_some())]);
    410     Ok(digest.finalize().into())
    411 }
    412 
    413 fn manifest_observation(
    414     source_id: RadrootsTradeEvidenceSourceIdV1,
    415     part: &RhiReconciliationReplayCommitParts,
    416     fact: &RhiReconciliationReplayCommitFact,
    417 ) -> Result<RadrootsTradeEvidenceManifestObservationV1, ()> {
    418     let record = &fact.record;
    419     let mut provenance = Sha256::new();
    420     provenance.update(PROVENANCE_DIGEST_DOMAIN);
    421     update_framed(&mut provenance, part.source_id.as_bytes())?;
    422     update_framed(&mut provenance, SOURCE_SELECTOR)?;
    423     provenance.update(part.policy_digest);
    424     provenance.update(record.event_id);
    425     provenance.update(record.event_signature);
    426     provenance.update(fact.observed_at.get().to_be_bytes());
    427     Ok(RadrootsTradeEvidenceManifestObservationV1::new(
    428         source_id,
    429         MutationId::from_bytes(record.mutation_id),
    430         EventId::from_bytes(record.event_id),
    431         RadrootsTradeSignedEventDigestV1::sha256(&record.canonical_event_json),
    432         RadrootsTradeEvidenceProvenanceDigestV1::from_bytes(provenance.finalize().into()),
    433     ))
    434 }
    435 
    436 fn update_framed(digest: &mut Sha256, bytes: &[u8]) -> Result<(), ()> {
    437     digest.update(u64::try_from(bytes.len()).map_err(|_| ())?.to_be_bytes());
    438     digest.update(bytes);
    439     Ok(())
    440 }
    441 
    442 fn update_optional_u64(digest: &mut Sha256, value: Option<u64>) {
    443     match value {
    444         Some(value) => {
    445             digest.update([1]);
    446             digest.update(value.to_be_bytes());
    447         }
    448         None => digest.update([0]),
    449     }
    450 }
    451 
    452 fn update_optional_cursor(digest: &mut Sha256, value: Option<crate::RhiTradeSourceCursor>) {
    453     match value {
    454         Some(value) => {
    455             digest.update([1]);
    456             digest.update(value.created_at_unix_seconds().to_be_bytes());
    457             digest.update(value.event_id());
    458         }
    459         None => digest.update([0]),
    460     }
    461 }
    462 
    463 const fn error(kind: RhiReconciliationManifestErrorKind) -> RhiReconciliationManifestError {
    464     RhiReconciliationManifestError { kind }
    465 }
    466 
    467 #[cfg(test)]
    468 mod tests {
    469     use super::*;
    470 
    471     #[test]
    472     fn diagnostics_are_closed_redacted_and_source_free() {
    473         for kind in [
    474             RhiReconciliationManifestErrorKind::InvalidObservationTime,
    475             RhiReconciliationManifestErrorKind::InvalidCommittedInventory,
    476         ] {
    477             let error = super::error(kind);
    478             assert_eq!(error.kind(), kind);
    479             assert!(Error::source(&error).is_none());
    480             assert!(!format!("{error} {error:?}").contains("trade-primary"));
    481         }
    482     }
    483 }