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