rhi

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

reconciliation_replay.rs (35820B)


      1 //! Pure overlap-safe reconciliation replay and provenance canonicalization.
      2 
      3 use core::{cmp::Ordering, fmt};
      4 use std::{collections::BTreeMap, error::Error};
      5 
      6 use sha2::{Digest, Sha256};
      7 
      8 use crate::{
      9     RhiAdmittedTradeMutationEvent, RhiConfigDocumentV1, RhiReconciliationAttemptPlan,
     10     RhiReconciliationSourceRequest, RhiReconciliationSourceRequestId,
     11     RhiReconciliationSourceResult, RhiReconciliationUnixMilliseconds,
     12     RhiTradeMutationObservedAtUnixSeconds, RhiTradeSourceCompletion, RhiTradeSourceCursor,
     13     state_metadata, state_trade::PersistenceRecord,
     14 };
     15 
     16 /// Exact version of the reconciliation replay contract.
     17 pub const RHI_RECONCILIATION_REPLAY_CONTRACT_VERSION: u32 = 1;
     18 
     19 const REPLAY_ID_DOMAIN: &[u8] = b"radroots.rhi.reconciliation_source_replay.v1\0";
     20 const SOURCE_KIND: &str = "nostr_relay";
     21 const SOURCE_SELECTOR: &str = "trade_mutation_lineage_v1";
     22 const MAX_OVERLAP_SECONDS: u64 = 86_400;
     23 
     24 /// Stable source-free replay validation class.
     25 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
     26 pub enum RhiReconciliationReplayErrorKind {
     27     InvalidInput,
     28     InvalidConfiguration,
     29     PolicyMismatch,
     30     ResourceLimit,
     31     MutationConflict,
     32     SignedEventConflict,
     33 }
     34 
     35 impl RhiReconciliationReplayErrorKind {
     36     /// Returns the stable machine-readable failure code.
     37     #[must_use]
     38     pub const fn code(self) -> &'static str {
     39         match self {
     40             Self::InvalidInput => "reconciliation_replay_input_invalid",
     41             Self::InvalidConfiguration => "reconciliation_replay_configuration_invalid",
     42             Self::PolicyMismatch => "reconciliation_replay_policy_mismatch",
     43             Self::ResourceLimit => "reconciliation_replay_resource_limit",
     44             Self::MutationConflict => "reconciliation_replay_mutation_conflict",
     45             Self::SignedEventConflict => "reconciliation_replay_signed_event_conflict",
     46         }
     47     }
     48 }
     49 
     50 /// Redacted source-free replay validation failure.
     51 #[derive(Clone, Copy, PartialEq, Eq)]
     52 pub struct RhiReconciliationReplayError {
     53     kind: RhiReconciliationReplayErrorKind,
     54 }
     55 
     56 impl RhiReconciliationReplayError {
     57     /// Returns the stable failure class.
     58     #[must_use]
     59     pub const fn kind(self) -> RhiReconciliationReplayErrorKind {
     60         self.kind
     61     }
     62 
     63     /// Returns the stable machine-readable failure code.
     64     #[must_use]
     65     pub const fn code(self) -> &'static str {
     66         self.kind.code()
     67     }
     68 }
     69 
     70 impl fmt::Display for RhiReconciliationReplayError {
     71     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
     72         formatter.write_str(match self.kind {
     73             RhiReconciliationReplayErrorKind::InvalidInput => {
     74                 "RHI reconciliation replay input is invalid"
     75             }
     76             RhiReconciliationReplayErrorKind::InvalidConfiguration => {
     77                 "RHI reconciliation replay configuration is invalid"
     78             }
     79             RhiReconciliationReplayErrorKind::PolicyMismatch => {
     80                 "RHI reconciliation replay policy does not match the attempt"
     81             }
     82             RhiReconciliationReplayErrorKind::ResourceLimit => {
     83                 "RHI reconciliation replay exceeds its resource limit"
     84             }
     85             RhiReconciliationReplayErrorKind::MutationConflict => {
     86                 "RHI reconciliation replay conflicts with canonical mutation evidence"
     87             }
     88             RhiReconciliationReplayErrorKind::SignedEventConflict => {
     89                 "RHI reconciliation replay conflicts with signed event evidence"
     90             }
     91         })
     92     }
     93 }
     94 
     95 impl fmt::Debug for RhiReconciliationReplayError {
     96     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
     97         formatter
     98             .debug_struct("RhiReconciliationReplayError")
     99             .field("kind", &self.kind)
    100             .finish()
    101     }
    102 }
    103 
    104 impl Error for RhiReconciliationReplayError {}
    105 
    106 /// Domain-separated identity of one request's exact cursor/overlap binding.
    107 #[derive(Clone, Copy, PartialEq, Eq, Hash)]
    108 pub struct RhiReconciliationSourceReplayId([u8; 32]);
    109 
    110 impl RhiReconciliationSourceReplayId {
    111     /// Returns the exact identity bytes.
    112     #[must_use]
    113     pub const fn as_bytes(&self) -> &[u8; 32] {
    114         &self.0
    115     }
    116 }
    117 
    118 impl fmt::Debug for RhiReconciliationSourceReplayId {
    119     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    120         formatter.write_str("RhiReconciliationSourceReplayId([redacted])")
    121     }
    122 }
    123 
    124 /// Sealed source/trade/policy/selector-scoped evidence for a committed cursor.
    125 ///
    126 /// Step 189 defines and consumes this non-forgeable capability but provides no
    127 /// minting path. Step 190 alone may construct it after the replay and cursor
    128 /// are committed atomically under their exact durable scope.
    129 #[derive(Clone, PartialEq, Eq)]
    130 pub struct RhiReconciliationSourceCursorEvidence {
    131     source_id: Box<str>,
    132     trade_id: radroots_event::id::TradeId,
    133     policy_digest: [u8; 32],
    134     selector_digest: [u8; 32],
    135     cursor: RhiTradeSourceCursor,
    136 }
    137 
    138 impl RhiReconciliationSourceCursorEvidence {
    139     /// Returns the retained exact cursor tuple.
    140     #[must_use]
    141     pub const fn cursor(&self) -> RhiTradeSourceCursor {
    142         self.cursor
    143     }
    144 }
    145 
    146 #[cfg(any(target_os = "linux", target_os = "macos"))]
    147 pub(crate) async fn read_committed_reconciliation_cursor(
    148     repositories: &crate::RhiStateRepositories<'_>,
    149     request: &crate::RhiReconciliationSourceRequest,
    150     policy: crate::RhiEvidencePolicyDigest,
    151 ) -> Result<Option<RhiReconciliationSourceCursorEvidence>, ()> {
    152     let source_id: Box<str> = request.source_id().into();
    153     let trade_id = request.trade_id();
    154     let selector_digest = *request.selector_digest().as_bytes();
    155     repositories
    156         .host()
    157         .sqlite_host()
    158         .transaction(move |transaction| {
    159             Box::pin(async move {
    160                 crate::source_ingest::read_checkpoint(
    161                     transaction,
    162                     source_id.as_ref(),
    163                     policy,
    164                     trade_id,
    165                 )
    166                 .await
    167                 .map(|checkpoint| {
    168                     checkpoint.map(|checkpoint| {
    169                         committed_cursor_evidence(
    170                             source_id,
    171                             trade_id,
    172                             *policy.as_bytes(),
    173                             selector_digest,
    174                             checkpoint.cursor,
    175                         )
    176                     })
    177                 })
    178             })
    179         })
    180         .await
    181         .map_err(|_| ())
    182 }
    183 
    184 impl fmt::Debug for RhiReconciliationSourceCursorEvidence {
    185     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    186         formatter
    187             .debug_struct("RhiReconciliationSourceCursorEvidence")
    188             .field("scope", &"[redacted]")
    189             .field("cursor", &self.cursor)
    190             .finish()
    191     }
    192 }
    193 
    194 /// Pure source-request binding to an optional prior cursor and overlap window.
    195 #[derive(Clone, PartialEq, Eq)]
    196 pub struct RhiReconciliationSourceReplayPlan {
    197     id: RhiReconciliationSourceReplayId,
    198     request_id: RhiReconciliationSourceRequestId,
    199     source_id: Box<str>,
    200     trade_id: radroots_event::id::TradeId,
    201     required: bool,
    202     policy_digest: [u8; 32],
    203     selector_digest: [u8; 32],
    204     prior_cursor: Option<RhiReconciliationSourceCursorEvidence>,
    205     overlap_seconds: u64,
    206     since_unix_seconds: u64,
    207 }
    208 
    209 impl RhiReconciliationSourceReplayPlan {
    210     /// Derives the exact overlap-safe cursor binding for one attempt request.
    211     pub fn from_request(
    212         attempt: &RhiReconciliationAttemptPlan,
    213         request: &RhiReconciliationSourceRequest,
    214         configuration: &RhiConfigDocumentV1,
    215         prior_cursor: Option<RhiReconciliationSourceCursorEvidence>,
    216     ) -> Result<Self, RhiReconciliationReplayError> {
    217         if !attempt
    218             .requests()
    219             .iter()
    220             .any(|candidate| candidate.id() == request.id())
    221         {
    222             return Err(failure(RhiReconciliationReplayErrorKind::InvalidInput));
    223         }
    224         let normalized = configuration.normalized();
    225         let policy = state_metadata::evidence_policy_digest(normalized)
    226             .map_err(|_| failure(RhiReconciliationReplayErrorKind::InvalidConfiguration))?;
    227         if policy != attempt.evidence_policy_digest() {
    228             return Err(failure(RhiReconciliationReplayErrorKind::PolicyMismatch));
    229         }
    230         if prior_cursor.as_ref().is_some_and(|evidence| {
    231             !cursor_scope_matches(
    232                 evidence,
    233                 request.source_id(),
    234                 request.trade_id(),
    235                 policy.as_bytes(),
    236                 request.selector_digest().as_bytes(),
    237             )
    238         }) {
    239             return Err(failure(RhiReconciliationReplayErrorKind::PolicyMismatch));
    240         }
    241         let source = normalized
    242             .pointer("/evidence/sources")
    243             .and_then(serde_json::Value::as_array)
    244             .and_then(|sources| {
    245                 sources.iter().find(|source| {
    246                     source
    247                         .pointer("/source_id")
    248                         .and_then(serde_json::Value::as_str)
    249                         == Some(request.source_id())
    250                 })
    251             })
    252             .filter(|source| {
    253                 source.pointer("/kind").and_then(serde_json::Value::as_str) == Some(SOURCE_KIND)
    254                     && source
    255                         .pointer("/selector")
    256                         .and_then(serde_json::Value::as_str)
    257                         == Some(SOURCE_SELECTOR)
    258             })
    259             .ok_or_else(|| failure(RhiReconciliationReplayErrorKind::InvalidConfiguration))?;
    260         let overlap_seconds = source
    261             .pointer("/overlap_seconds")
    262             .and_then(serde_json::Value::as_u64)
    263             .filter(|value| (1..=MAX_OVERLAP_SECONDS).contains(value))
    264             .ok_or_else(|| failure(RhiReconciliationReplayErrorKind::InvalidConfiguration))?;
    265         let lookback_seconds = source
    266             .pointer("/lookback_seconds")
    267             .and_then(serde_json::Value::as_u64)
    268             .filter(|value| *value == request.lookback_seconds() && overlap_seconds <= *value)
    269             .ok_or_else(|| failure(RhiReconciliationReplayErrorKind::InvalidConfiguration))?;
    270         let raw_prior_cursor = prior_cursor.as_ref().map(|evidence| evidence.cursor);
    271         let since_unix_seconds = raw_prior_cursor.map_or_else(
    272             || {
    273                 request
    274                     .attempt_started_at()
    275                     .get()
    276                     .checked_div(1_000)
    277                     .unwrap_or(0)
    278                     .saturating_sub(lookback_seconds)
    279             },
    280             |cursor| resume_since(cursor, overlap_seconds),
    281         );
    282         Ok(Self {
    283             id: replay_id(
    284                 request.id(),
    285                 raw_prior_cursor,
    286                 overlap_seconds,
    287                 since_unix_seconds,
    288             ),
    289             request_id: request.id(),
    290             source_id: request.source_id().into(),
    291             trade_id: request.trade_id(),
    292             required: request.required(),
    293             policy_digest: *policy.as_bytes(),
    294             selector_digest: *request.selector_digest().as_bytes(),
    295             prior_cursor,
    296             overlap_seconds,
    297             since_unix_seconds,
    298         })
    299     }
    300 
    301     /// Returns the exact cursor-binding identity.
    302     #[must_use]
    303     pub const fn id(&self) -> RhiReconciliationSourceReplayId {
    304         self.id
    305     }
    306 
    307     /// Returns the exact source request identity.
    308     #[must_use]
    309     pub const fn request_id(&self) -> RhiReconciliationSourceRequestId {
    310         self.request_id
    311     }
    312 
    313     /// Returns the retained prior cursor, when one exists.
    314     #[must_use]
    315     pub fn prior_cursor(&self) -> Option<RhiTradeSourceCursor> {
    316         self.prior_cursor.as_ref().map(|evidence| evidence.cursor)
    317     }
    318 
    319     /// Returns the configured overlap in whole seconds.
    320     #[must_use]
    321     pub const fn overlap_seconds(&self) -> u64 {
    322         self.overlap_seconds
    323     }
    324 
    325     /// Returns the inclusive overlap-safe query start in Unix seconds.
    326     #[must_use]
    327     pub const fn since_unix_seconds(&self) -> u64 {
    328         self.since_unix_seconds
    329     }
    330 
    331     /// Canonicalizes a bounded admitted-event result for a later atomic commit.
    332     pub fn finish<I>(
    333         self,
    334         request: &RhiReconciliationSourceRequest,
    335         outcome: RhiTradeSourceCompletion,
    336         started_at: RhiReconciliationUnixMilliseconds,
    337         finished_at: RhiReconciliationUnixMilliseconds,
    338         events: I,
    339     ) -> Result<RhiReconciliationSourceReplay, RhiReconciliationReplayError>
    340     where
    341         I: IntoIterator<Item = RhiAdmittedTradeMutationEvent>,
    342     {
    343         if request.id() != self.request_id {
    344             return Err(failure(RhiReconciliationReplayErrorKind::InvalidInput));
    345         }
    346         RhiReconciliationSourceResult::new(request, outcome, started_at, finished_at, 0, 0)
    347             .map_err(|_| failure(RhiReconciliationReplayErrorKind::InvalidInput))?;
    348         let maximum_events = usize::try_from(request.maximum_events())
    349             .map_err(|_| failure(RhiReconciliationReplayErrorKind::InvalidConfiguration))?;
    350         let candidates = ingest_bounded_facts(
    351             events.into_iter().map(|event| {
    352                 if event.mutation().trade_id != request.trade_id() {
    353                     return Err(failure(RhiReconciliationReplayErrorKind::InvalidInput));
    354                 }
    355                 let event_bytes = u64::try_from(event.original_bytes().len())
    356                     .map_err(|_| failure(RhiReconciliationReplayErrorKind::ResourceLimit))?;
    357                 let observed_at = event.observed_at_unix_seconds();
    358                 let observed_at_milliseconds = observed_at
    359                     .get()
    360                     .checked_mul(1_000)
    361                     .ok_or_else(|| failure(RhiReconciliationReplayErrorKind::InvalidInput))?;
    362                 let observed_interval_end = observed_at_milliseconds
    363                     .checked_add(999)
    364                     .ok_or_else(|| failure(RhiReconciliationReplayErrorKind::InvalidInput))?;
    365                 if observed_interval_end < started_at.get()
    366                     || observed_at_milliseconds > finished_at.get()
    367                 {
    368                     return Err(failure(RhiReconciliationReplayErrorKind::InvalidInput));
    369                 }
    370                 Ok(ReplayFact {
    371                     original_bytes: event_bytes,
    372                     observed_at,
    373                     record: PersistenceRecord::from_admitted(event)
    374                         .map_err(|_| failure(RhiReconciliationReplayErrorKind::InvalidInput))?,
    375                 })
    376             }),
    377             maximum_events,
    378             request.maximum_bytes(),
    379         )?;
    380         let CanonicalReplay {
    381             facts,
    382             duplicate_observations,
    383             accepted_original_bytes,
    384             cursor_candidate,
    385             first_observed_at,
    386         } = canonicalize(candidates)?;
    387         let accepted_event_count = u32::try_from(facts.len())
    388             .map_err(|_| failure(RhiReconciliationReplayErrorKind::ResourceLimit))?;
    389         let result = RhiReconciliationSourceResult::new(
    390             request,
    391             outcome,
    392             started_at,
    393             finished_at,
    394             accepted_event_count,
    395             accepted_original_bytes,
    396         )
    397         .map_err(|_| failure(RhiReconciliationReplayErrorKind::InvalidInput))?;
    398         Ok(RhiReconciliationSourceReplay {
    399             plan: self,
    400             result,
    401             duplicate_observations,
    402             accepted_original_bytes,
    403             cursor_candidate,
    404             first_observed_at,
    405             facts: facts.into_boxed_slice(),
    406         })
    407     }
    408 }
    409 
    410 impl fmt::Debug for RhiReconciliationSourceReplayPlan {
    411     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    412         formatter
    413             .debug_struct("RhiReconciliationSourceReplayPlan")
    414             .field("identity", &"[redacted]")
    415             .field("has_prior_cursor", &self.prior_cursor.is_some())
    416             .field("overlap_seconds", &self.overlap_seconds)
    417             .field("since_unix_seconds", &self.since_unix_seconds)
    418             .finish()
    419     }
    420 }
    421 
    422 /// Canonical bounded replay result prepared for the Step 190 commit boundary.
    423 pub struct RhiReconciliationSourceReplay {
    424     plan: RhiReconciliationSourceReplayPlan,
    425     result: RhiReconciliationSourceResult,
    426     duplicate_observations: u32,
    427     accepted_original_bytes: u64,
    428     cursor_candidate: Option<RhiTradeSourceCursor>,
    429     first_observed_at: Option<RhiTradeMutationObservedAtUnixSeconds>,
    430     facts: Box<[ReplayFact]>,
    431 }
    432 
    433 impl RhiReconciliationSourceReplay {
    434     /// Returns the exact cursor-binding identity.
    435     #[must_use]
    436     pub const fn id(&self) -> RhiReconciliationSourceReplayId {
    437         self.plan.id()
    438     }
    439 
    440     /// Returns the validated source result derived from the canonical inventory.
    441     #[must_use]
    442     pub const fn result(&self) -> RhiReconciliationSourceResult {
    443         self.result
    444     }
    445 
    446     /// Returns the number of distinct accepted signed-event identities.
    447     #[must_use]
    448     pub fn accepted_event_count(&self) -> usize {
    449         self.facts.len()
    450     }
    451 
    452     /// Returns the retained first-provenance original-byte total.
    453     #[must_use]
    454     pub const fn accepted_original_event_bytes(&self) -> u64 {
    455         self.accepted_original_bytes
    456     }
    457 
    458     /// Returns the exact replay observations removed by canonical deduplication.
    459     #[must_use]
    460     pub const fn duplicate_observation_count(&self) -> u32 {
    461         self.duplicate_observations
    462     }
    463 
    464     /// Returns the earliest retained source observation, when evidence exists.
    465     #[must_use]
    466     pub const fn first_observed_at(&self) -> Option<RhiTradeMutationObservedAtUnixSeconds> {
    467         self.first_observed_at
    468     }
    469 
    470     /// Returns the greatest admitted cursor candidate regardless of completion.
    471     #[must_use]
    472     pub const fn cursor_candidate(&self) -> Option<RhiTradeSourceCursor> {
    473         self.cursor_candidate
    474     }
    475 
    476     /// Returns a cursor only when exact source completion makes it eligible.
    477     #[must_use]
    478     pub fn eligible_cursor(&self) -> Option<RhiTradeSourceCursor> {
    479         eligible_cursor(
    480             self.result.outcome(),
    481             self.cursor_candidate,
    482             self.plan
    483                 .prior_cursor
    484                 .as_ref()
    485                 .map(|evidence| evidence.cursor),
    486         )
    487     }
    488 
    489     pub(crate) fn into_commit_parts(self) -> RhiReconciliationReplayCommitParts {
    490         let eligible_cursor = self.eligible_cursor();
    491         RhiReconciliationReplayCommitParts {
    492             replay_id: self.plan.id,
    493             request_id: self.plan.request_id,
    494             source_id: self.plan.source_id,
    495             trade_id: self.plan.trade_id,
    496             required: self.plan.required,
    497             policy_digest: self.plan.policy_digest,
    498             selector_digest: self.plan.selector_digest,
    499             prior_cursor: self.plan.prior_cursor.map(|evidence| evidence.cursor),
    500             overlap_seconds: self.plan.overlap_seconds,
    501             since_unix_seconds: self.plan.since_unix_seconds,
    502             result: self.result,
    503             duplicate_observations: self.duplicate_observations,
    504             cursor_candidate: self.cursor_candidate,
    505             eligible_cursor,
    506             first_observed_at: self.first_observed_at,
    507             facts: self
    508                 .facts
    509                 .into_vec()
    510                 .into_iter()
    511                 .map(|fact| RhiReconciliationReplayCommitFact {
    512                     record: fact.record,
    513                     observed_at: fact.observed_at,
    514                 })
    515                 .collect::<Vec<_>>()
    516                 .into_boxed_slice(),
    517         }
    518     }
    519 }
    520 
    521 pub(crate) struct RhiReconciliationReplayCommitParts {
    522     pub(crate) replay_id: RhiReconciliationSourceReplayId,
    523     pub(crate) request_id: RhiReconciliationSourceRequestId,
    524     pub(crate) source_id: Box<str>,
    525     pub(crate) trade_id: radroots_event::id::TradeId,
    526     pub(crate) required: bool,
    527     pub(crate) policy_digest: [u8; 32],
    528     pub(crate) selector_digest: [u8; 32],
    529     pub(crate) prior_cursor: Option<RhiTradeSourceCursor>,
    530     pub(crate) overlap_seconds: u64,
    531     pub(crate) since_unix_seconds: u64,
    532     pub(crate) result: RhiReconciliationSourceResult,
    533     pub(crate) duplicate_observations: u32,
    534     pub(crate) cursor_candidate: Option<RhiTradeSourceCursor>,
    535     pub(crate) eligible_cursor: Option<RhiTradeSourceCursor>,
    536     pub(crate) first_observed_at: Option<RhiTradeMutationObservedAtUnixSeconds>,
    537     pub(crate) facts: Box<[RhiReconciliationReplayCommitFact]>,
    538 }
    539 
    540 pub(crate) struct RhiReconciliationReplayCommitFact {
    541     pub(crate) record: PersistenceRecord,
    542     pub(crate) observed_at: RhiTradeMutationObservedAtUnixSeconds,
    543 }
    544 
    545 pub(crate) fn committed_cursor_evidence(
    546     source_id: Box<str>,
    547     trade_id: radroots_event::id::TradeId,
    548     policy_digest: [u8; 32],
    549     selector_digest: [u8; 32],
    550     cursor: RhiTradeSourceCursor,
    551 ) -> RhiReconciliationSourceCursorEvidence {
    552     RhiReconciliationSourceCursorEvidence {
    553         source_id,
    554         trade_id,
    555         policy_digest,
    556         selector_digest,
    557         cursor,
    558     }
    559 }
    560 
    561 fn cursor_scope_matches(
    562     evidence: &RhiReconciliationSourceCursorEvidence,
    563     source_id: &str,
    564     trade_id: radroots_event::id::TradeId,
    565     policy_digest: &[u8; 32],
    566     selector_digest: &[u8; 32],
    567 ) -> bool {
    568     evidence.source_id.as_ref() == source_id
    569         && evidence.trade_id == trade_id
    570         && &evidence.policy_digest == policy_digest
    571         && &evidence.selector_digest == selector_digest
    572 }
    573 
    574 impl fmt::Debug for RhiReconciliationSourceReplay {
    575     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    576         formatter
    577             .debug_struct("RhiReconciliationSourceReplay")
    578             .field("identity", &"[redacted]")
    579             .field("outcome", &self.result.outcome())
    580             .field("accepted_event_count", &self.facts.len())
    581             .field(
    582                 "accepted_original_event_bytes",
    583                 &self.accepted_original_bytes,
    584             )
    585             .field("duplicate_observations", &self.duplicate_observations)
    586             .field("has_cursor_candidate", &self.cursor_candidate.is_some())
    587             .finish()
    588     }
    589 }
    590 
    591 struct ReplayFact {
    592     record: PersistenceRecord,
    593     original_bytes: u64,
    594     observed_at: RhiTradeMutationObservedAtUnixSeconds,
    595 }
    596 
    597 struct CanonicalReplay {
    598     facts: Vec<ReplayFact>,
    599     duplicate_observations: u32,
    600     accepted_original_bytes: u64,
    601     cursor_candidate: Option<RhiTradeSourceCursor>,
    602     first_observed_at: Option<RhiTradeMutationObservedAtUnixSeconds>,
    603 }
    604 
    605 fn ingest_bounded_facts<I>(
    606     facts: I,
    607     maximum_events: usize,
    608     maximum_bytes: u64,
    609 ) -> Result<Vec<ReplayFact>, RhiReconciliationReplayError>
    610 where
    611     I: IntoIterator<Item = Result<ReplayFact, RhiReconciliationReplayError>>,
    612 {
    613     let mut accepted_bytes = 0_u64;
    614     let mut bounded = Vec::with_capacity(maximum_events);
    615     for fact in facts.into_iter().take(maximum_events.saturating_add(1)) {
    616         if bounded.len() == maximum_events {
    617             return Err(failure(RhiReconciliationReplayErrorKind::ResourceLimit));
    618         }
    619         let fact = fact?;
    620         accepted_bytes = accepted_bytes
    621             .checked_add(fact.original_bytes)
    622             .filter(|value| *value <= maximum_bytes)
    623             .ok_or_else(|| failure(RhiReconciliationReplayErrorKind::ResourceLimit))?;
    624         bounded.push(fact);
    625     }
    626     Ok(bounded)
    627 }
    628 
    629 fn canonicalize(
    630     mut candidates: Vec<ReplayFact>,
    631 ) -> Result<CanonicalReplay, RhiReconciliationReplayError> {
    632     candidates.sort_by(compare_fact);
    633     let mut event_indexes = BTreeMap::<([u8; 32], [u8; 64]), usize>::new();
    634     let mut mutation_indexes = BTreeMap::<[u8; 32], usize>::new();
    635     let mut facts = Vec::<ReplayFact>::with_capacity(candidates.len());
    636     let mut duplicate_observations = 0_u32;
    637     let mut accepted_original_bytes = 0_u64;
    638     let mut cursor_candidate = None;
    639     let mut first_observed_at: Option<RhiTradeMutationObservedAtUnixSeconds> = None;
    640     for candidate in candidates {
    641         let event_key = (candidate.record.event_id, candidate.record.event_signature);
    642         if let Some(index) = event_indexes.get(&event_key).copied() {
    643             if !same_event(&facts[index].record, &candidate.record) {
    644                 return Err(failure(
    645                     RhiReconciliationReplayErrorKind::SignedEventConflict,
    646                 ));
    647             }
    648             duplicate_observations = duplicate_observations
    649                 .checked_add(1)
    650                 .ok_or_else(|| failure(RhiReconciliationReplayErrorKind::ResourceLimit))?;
    651             continue;
    652         }
    653         if let Some(index) = mutation_indexes.get(&candidate.record.mutation_id).copied()
    654             && !same_mutation(&facts[index].record, &candidate.record)
    655         {
    656             return Err(failure(RhiReconciliationReplayErrorKind::MutationConflict));
    657         }
    658         accepted_original_bytes = accepted_original_bytes
    659             .checked_add(candidate.original_bytes)
    660             .ok_or_else(|| failure(RhiReconciliationReplayErrorKind::ResourceLimit))?;
    661         let cursor = RhiTradeSourceCursor::from_verified_parts(
    662             candidate.record.authored_at_unix_s,
    663             candidate.record.event_id,
    664         );
    665         cursor_candidate = Some(cursor_candidate.map_or(cursor, |current| {
    666             if compare_cursor(current, cursor).is_lt() {
    667                 cursor
    668             } else {
    669                 current
    670             }
    671         }));
    672         first_observed_at = Some(first_observed_at.map_or(candidate.observed_at, |current| {
    673             if candidate.observed_at.get() < current.get() {
    674                 candidate.observed_at
    675             } else {
    676                 current
    677             }
    678         }));
    679         let index = facts.len();
    680         event_indexes.insert(event_key, index);
    681         mutation_indexes
    682             .entry(candidate.record.mutation_id)
    683             .or_insert(index);
    684         facts.push(candidate);
    685     }
    686     Ok(CanonicalReplay {
    687         facts,
    688         duplicate_observations,
    689         accepted_original_bytes,
    690         cursor_candidate,
    691         first_observed_at,
    692     })
    693 }
    694 
    695 fn compare_fact(left: &ReplayFact, right: &ReplayFact) -> Ordering {
    696     left.record
    697         .authored_at_unix_s
    698         .cmp(&right.record.authored_at_unix_s)
    699         .then_with(|| left.record.event_id.cmp(&right.record.event_id))
    700         .then_with(|| {
    701             left.record
    702                 .event_signature
    703                 .cmp(&right.record.event_signature)
    704         })
    705         .then_with(|| left.observed_at.get().cmp(&right.observed_at.get()))
    706         .then_with(|| left.original_bytes.cmp(&right.original_bytes))
    707 }
    708 
    709 fn same_mutation(left: &PersistenceRecord, right: &PersistenceRecord) -> bool {
    710     left.mutation_id == right.mutation_id
    711         && left.trade_id == right.trade_id
    712         && left.contract_id == right.contract_id
    713         && left.schema_version == right.schema_version
    714         && left.event_kind == right.event_kind
    715         && left.author_pubkey == right.author_pubkey
    716         && left.canonical_content == right.canonical_content
    717 }
    718 
    719 fn same_event(left: &PersistenceRecord, right: &PersistenceRecord) -> bool {
    720     same_mutation(left, right)
    721         && left.event_id == right.event_id
    722         && left.event_signature == right.event_signature
    723         && left.event_kind == right.event_kind
    724         && left.authored_at_unix_s == right.authored_at_unix_s
    725         && left.canonical_event_json == right.canonical_event_json
    726 }
    727 
    728 fn compare_cursor(left: RhiTradeSourceCursor, right: RhiTradeSourceCursor) -> Ordering {
    729     (left.created_at_unix_seconds(), left.event_id())
    730         .cmp(&(right.created_at_unix_seconds(), right.event_id()))
    731 }
    732 
    733 fn resume_since(cursor: RhiTradeSourceCursor, overlap_seconds: u64) -> u64 {
    734     cursor
    735         .created_at_unix_seconds()
    736         .saturating_sub(overlap_seconds)
    737 }
    738 
    739 fn eligible_cursor(
    740     outcome: RhiTradeSourceCompletion,
    741     candidate: Option<RhiTradeSourceCursor>,
    742     prior: Option<RhiTradeSourceCursor>,
    743 ) -> Option<RhiTradeSourceCursor> {
    744     if !matches!(outcome, RhiTradeSourceCompletion::Complete) {
    745         return None;
    746     }
    747     candidate
    748         .filter(|candidate| prior.is_none_or(|prior| compare_cursor(prior, *candidate).is_lt()))
    749 }
    750 
    751 fn replay_id(
    752     request_id: RhiReconciliationSourceRequestId,
    753     prior_cursor: Option<RhiTradeSourceCursor>,
    754     overlap_seconds: u64,
    755     since_unix_seconds: u64,
    756 ) -> RhiReconciliationSourceReplayId {
    757     let mut digest = Sha256::new();
    758     digest.update(REPLAY_ID_DOMAIN);
    759     digest.update(request_id.as_bytes());
    760     match prior_cursor {
    761         Some(cursor) => {
    762             digest.update([1]);
    763             digest.update(cursor.created_at_unix_seconds().to_be_bytes());
    764             digest.update(cursor.event_id());
    765         }
    766         None => digest.update([0]),
    767     }
    768     digest.update(overlap_seconds.to_be_bytes());
    769     digest.update(since_unix_seconds.to_be_bytes());
    770     RhiReconciliationSourceReplayId(digest.finalize().into())
    771 }
    772 
    773 const fn failure(kind: RhiReconciliationReplayErrorKind) -> RhiReconciliationReplayError {
    774     RhiReconciliationReplayError { kind }
    775 }
    776 
    777 #[cfg(test)]
    778 mod tests {
    779     use super::*;
    780 
    781     fn observed(value: u64) -> RhiTradeMutationObservedAtUnixSeconds {
    782         RhiTradeMutationObservedAtUnixSeconds::new(value).expect("observation")
    783     }
    784 
    785     fn record(event: u8, signature: u8, mutation: u8, content: &[u8]) -> PersistenceRecord {
    786         PersistenceRecord {
    787             mutation_id: [mutation; 32],
    788             trade_id: [0x11; 16],
    789             contract_id: "radroots.trade.proposal.v1",
    790             schema_version: 1,
    791             event_id: [event; 32],
    792             event_signature: [signature; 64],
    793             author_pubkey: [0x22; 32],
    794             event_kind: 3470,
    795             authored_at_unix_s: 1_784_347_200,
    796             canonical_content: content.into(),
    797             canonical_event_json: [b"event:".as_slice(), content].concat().into_boxed_slice(),
    798         }
    799     }
    800 
    801     fn fact(
    802         event: u8,
    803         signature: u8,
    804         mutation: u8,
    805         content: &[u8],
    806         observation: u64,
    807         bytes: u64,
    808     ) -> ReplayFact {
    809         ReplayFact {
    810             record: record(event, signature, mutation, content),
    811             original_bytes: bytes,
    812             observed_at: observed(observation),
    813         }
    814     }
    815 
    816     #[test]
    817     fn canonicalization_deduplicates_replay_and_retains_first_provenance() {
    818         let canonical = canonicalize(vec![
    819             fact(1, 2, 3, b"same", 1_784_347_202, 12),
    820             fact(1, 2, 3, b"same", 1_784_347_200, 10),
    821             fact(1, 4, 3, b"same", 1_784_347_201, 11),
    822         ])
    823         .expect("canonical replay");
    824         assert_eq!(canonical.facts.len(), 2);
    825         assert_eq!(canonical.duplicate_observations, 1);
    826         assert_eq!(canonical.accepted_original_bytes, 21);
    827         assert_eq!(
    828             canonical.first_observed_at.expect("first").get(),
    829             1_784_347_200
    830         );
    831     }
    832 
    833     #[test]
    834     fn conflicting_mutation_and_event_identity_reuse_fail_closed() {
    835         let mutation = canonicalize(vec![
    836             fact(1, 2, 3, b"first", 1_784_347_200, 10),
    837             fact(4, 5, 3, b"second", 1_784_347_201, 10),
    838         ])
    839         .err()
    840         .expect("mutation conflict");
    841         assert_eq!(
    842             mutation.kind(),
    843             RhiReconciliationReplayErrorKind::MutationConflict
    844         );
    845 
    846         let event = canonicalize(vec![
    847             fact(1, 2, 3, b"first", 1_784_347_200, 10),
    848             fact(1, 2, 4, b"second", 1_784_347_201, 10),
    849         ])
    850         .err()
    851         .expect("event conflict");
    852         assert_eq!(
    853             event.kind(),
    854             RhiReconciliationReplayErrorKind::SignedEventConflict
    855         );
    856     }
    857 
    858     #[test]
    859     fn pre_dedup_bound_is_exact_and_infinite_iterators_terminate() {
    860         let exact = ingest_bounded_facts([Ok(fact(1, 2, 3, b"one", 1_784_347_200, 10))], 1, 10)
    861             .expect("exact maximum");
    862         assert_eq!(exact.len(), 1);
    863 
    864         let over_bytes =
    865             ingest_bounded_facts([Ok(fact(1, 2, 3, b"one", 1_784_347_200, 11))], 1, 10)
    866                 .err()
    867                 .expect("byte maximum plus one");
    868         assert_eq!(
    869             over_bytes.kind(),
    870             RhiReconciliationReplayErrorKind::ResourceLimit
    871         );
    872 
    873         let mut event = 0_u8;
    874         let over_count = ingest_bounded_facts(
    875             std::iter::repeat_with(|| {
    876                 event = event.wrapping_add(1);
    877                 Ok(fact(event, 2, event, b"one", 1_784_347_200, 1))
    878             }),
    879             1,
    880             10,
    881         )
    882         .err()
    883         .expect("count maximum plus one");
    884         assert_eq!(
    885             over_count.kind(),
    886             RhiReconciliationReplayErrorKind::ResourceLimit
    887         );
    888     }
    889 
    890     #[test]
    891     fn errors_are_source_free_and_redacted() {
    892         let error = failure(RhiReconciliationReplayErrorKind::SignedEventConflict);
    893         assert!(Error::source(&error).is_none());
    894         assert_eq!(error.code(), "reconciliation_replay_signed_event_conflict");
    895         assert_eq!(
    896             format!("{error:?}"),
    897             "RhiReconciliationReplayError { kind: SignedEventConflict }"
    898         );
    899     }
    900 
    901     #[test]
    902     fn cursor_evidence_scope_requires_every_exact_dimension() {
    903         let exact = RhiReconciliationSourceCursorEvidence {
    904             source_id: "trade-primary".into(),
    905             trade_id: radroots_event::id::TradeId::from_bytes([0x11; 16]),
    906             policy_digest: [0x22; 32],
    907             selector_digest: [0x33; 32],
    908             cursor: RhiTradeSourceCursor::from_verified_parts(42, [0x44; 32]),
    909         };
    910         assert!(cursor_scope_matches(
    911             &exact,
    912             "trade-primary",
    913             radroots_event::id::TradeId::from_bytes([0x11; 16]),
    914             &[0x22; 32],
    915             &[0x33; 32],
    916         ));
    917         assert!(!cursor_scope_matches(
    918             &exact,
    919             "trade-secondary",
    920             radroots_event::id::TradeId::from_bytes([0x11; 16]),
    921             &[0x22; 32],
    922             &[0x33; 32],
    923         ));
    924         assert!(!cursor_scope_matches(
    925             &exact,
    926             "trade-primary",
    927             radroots_event::id::TradeId::from_bytes([0x12; 16]),
    928             &[0x22; 32],
    929             &[0x33; 32],
    930         ));
    931         assert!(!cursor_scope_matches(
    932             &exact,
    933             "trade-primary",
    934             radroots_event::id::TradeId::from_bytes([0x11; 16]),
    935             &[0x23; 32],
    936             &[0x33; 32],
    937         ));
    938         assert!(!cursor_scope_matches(
    939             &exact,
    940             "trade-primary",
    941             radroots_event::id::TradeId::from_bytes([0x11; 16]),
    942             &[0x22; 32],
    943             &[0x34; 32],
    944         ));
    945     }
    946 
    947     #[test]
    948     fn resume_overlap_and_cursor_eligibility_are_exact() {
    949         let prior = RhiTradeSourceCursor::from_verified_parts(500, [0x11; 32]);
    950         let equal = RhiTradeSourceCursor::from_verified_parts(500, [0x11; 32]);
    951         let older = RhiTradeSourceCursor::from_verified_parts(499, [0xff; 32]);
    952         let newer = RhiTradeSourceCursor::from_verified_parts(500, [0x12; 32]);
    953         assert_eq!(resume_since(prior, 300), 200);
    954         assert_eq!(resume_since(prior, 600), 0);
    955         assert_eq!(
    956             eligible_cursor(RhiTradeSourceCompletion::Complete, Some(newer), Some(prior)),
    957             Some(newer)
    958         );
    959         assert_eq!(
    960             eligible_cursor(RhiTradeSourceCompletion::Complete, Some(equal), Some(prior)),
    961             None
    962         );
    963         assert_eq!(
    964             eligible_cursor(RhiTradeSourceCompletion::Complete, Some(older), Some(prior)),
    965             None
    966         );
    967         assert_eq!(
    968             eligible_cursor(
    969                 RhiTradeSourceCompletion::IncompleteUnavailable,
    970                 Some(newer),
    971                 Some(prior),
    972             ),
    973             None
    974         );
    975     }
    976 }