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 }