trade_ingest.rs (33375B)
1 //! Allocation-bounded, cryptographically verified trade-mutation admission. 2 3 use core::fmt; 4 use std::{borrow::Cow, collections::BTreeSet, error::Error}; 5 6 use radroots_event::{ 7 envelope::{EventEnvelope, kind::is_trade_mutation_event_kind}, 8 id::{EventId, MutationId}, 9 trade::TradeMutationEnvelopeV1, 10 wire::{ 11 DEFAULT_EXTRA_MAX_FIELDS, DEFAULT_EXTRA_TOTAL_JSON_MAX_BYTES, EventWireLimits, 12 Nip01EventWire, 13 }, 14 }; 15 use radroots_event_codec::decode::trade::{RadrootsTradeMutationError, trade_mutation_from_event}; 16 use radroots_nostr::event::{Verification, verify, verify_id}; 17 use serde::Deserialize; 18 use serde::de::{self, DeserializeSeed, IgnoredAny, MapAccess, SeqAccess, Visitor}; 19 use serde_json::value::RawValue; 20 21 use crate::RhiConfigDocumentV1; 22 23 /// Maximum encoded Nostr event-identifier length admitted before allocation. 24 pub const RHI_TRADE_EVENT_ID_MAX_BYTES: usize = 64; 25 26 /// Exact version of the RHI trade-ingest contract. 27 pub const RHI_TRADE_INGEST_CONTRACT_VERSION: u32 = 1; 28 29 /// Maximum encoded Nostr public-key length admitted before allocation. 30 pub const RHI_TRADE_EVENT_PUBLIC_KEY_MAX_BYTES: usize = 64; 31 32 /// Maximum encoded Nostr signature length admitted before allocation. 33 pub const RHI_TRADE_EVENT_SIGNATURE_MAX_BYTES: usize = 128; 34 35 /// Maximum number of bounded, non-authoritative outer event extensions. 36 pub const RHI_TRADE_EVENT_EXTRA_FIELD_MAX_COUNT: usize = DEFAULT_EXTRA_MAX_FIELDS; 37 38 /// Maximum aggregate JSON bytes for non-authoritative outer event extensions. 39 pub const RHI_TRADE_EVENT_EXTRA_JSON_MAX_BYTES: usize = DEFAULT_EXTRA_TOTAL_JSON_MAX_BYTES; 40 41 const DUPLICATE_FIELD_SENTINEL: &str = "rhi-duplicate-event-field"; 42 const EXTRA_COUNT_SENTINEL: &str = "rhi-extra-field-count-limit"; 43 const EXTRA_BYTES_SENTINEL: &str = "rhi-extra-field-bytes-limit"; 44 const TAG_COUNT_SENTINEL: &str = "rhi-tag-count-limit"; 45 const TAG_ELEMENT_COUNT_SENTINEL: &str = "rhi-tag-element-count-limit"; 46 const TAG_ELEMENT_BYTES_SENTINEL: &str = "rhi-tag-element-bytes-limit"; 47 const TAG_TOTAL_BYTES_SENTINEL: &str = "rhi-tag-total-bytes-limit"; 48 49 /// Immutable trade-event limits projected from one admitted RHI configuration. 50 #[derive(Clone, Copy, PartialEq, Eq)] 51 pub struct RhiTradeMutationAdmissionLimits { 52 wire_bytes: usize, 53 content_bytes: usize, 54 tag_count: usize, 55 tag_total_elements: usize, 56 tag_element_bytes: usize, 57 tag_total_bytes: usize, 58 } 59 60 impl RhiTradeMutationAdmissionLimits { 61 /// Projects the exact event limits from a validated immutable configuration. 62 pub fn from_config( 63 configuration: &RhiConfigDocumentV1, 64 ) -> Result<Self, RhiTradeMutationAdmissionError> { 65 Ok(Self { 66 wire_bytes: config_limit(configuration, "/resource_limits/events/wire_bytes")?, 67 content_bytes: config_limit(configuration, "/resource_limits/events/content_bytes")?, 68 tag_count: config_limit(configuration, "/resource_limits/events/tag_count")?, 69 tag_total_elements: config_limit( 70 configuration, 71 "/resource_limits/events/tag_total_elements", 72 )?, 73 tag_element_bytes: config_limit( 74 configuration, 75 "/resource_limits/events/tag_element_bytes", 76 )?, 77 tag_total_bytes: config_limit( 78 configuration, 79 "/resource_limits/events/tag_total_bytes", 80 )?, 81 }) 82 } 83 84 /// Returns the original event-wire byte cap. 85 #[must_use] 86 pub const fn wire_bytes(self) -> usize { 87 self.wire_bytes 88 } 89 90 /// Returns the decoded canonical-content byte cap. 91 #[must_use] 92 pub const fn content_bytes(self) -> usize { 93 self.content_bytes 94 } 95 96 /// Returns the event-tag count cap. 97 #[must_use] 98 pub const fn tag_count(self) -> usize { 99 self.tag_count 100 } 101 102 /// Returns the aggregate event-tag-element count cap. 103 #[must_use] 104 pub const fn tag_total_elements(self) -> usize { 105 self.tag_total_elements 106 } 107 108 /// Returns the decoded byte cap for one tag element. 109 #[must_use] 110 pub const fn tag_element_bytes(self) -> usize { 111 self.tag_element_bytes 112 } 113 114 /// Returns the aggregate decoded byte cap for all tag elements. 115 #[must_use] 116 pub const fn tag_total_bytes(self) -> usize { 117 self.tag_total_bytes 118 } 119 120 const fn wire_limits(self) -> EventWireLimits { 121 EventWireLimits { 122 max_raw_json_bytes: self.wire_bytes, 123 max_content_bytes: self.content_bytes, 124 max_tag_count: self.tag_count, 125 max_total_tag_elements: self.tag_total_elements, 126 max_tag_element_bytes: self.tag_element_bytes, 127 max_total_tag_bytes: self.tag_total_bytes, 128 max_extra_fields: RHI_TRADE_EVENT_EXTRA_FIELD_MAX_COUNT, 129 max_total_extra_json_bytes: RHI_TRADE_EVENT_EXTRA_JSON_MAX_BYTES, 130 } 131 } 132 } 133 134 impl fmt::Debug for RhiTradeMutationAdmissionLimits { 135 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 136 formatter 137 .debug_struct("RhiTradeMutationAdmissionLimits") 138 .field("wire_bytes", &self.wire_bytes) 139 .field("content_bytes", &self.content_bytes) 140 .field("tag_count", &self.tag_count) 141 .field("tag_total_elements", &self.tag_total_elements) 142 .field("tag_element_bytes", &self.tag_element_bytes) 143 .field("tag_total_bytes", &self.tag_total_bytes) 144 .finish() 145 } 146 } 147 148 /// Injected UTC second at which one trade event is observed. 149 #[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Hash)] 150 pub struct RhiTradeMutationObservedAtUnixSeconds(u64); 151 152 impl RhiTradeMutationObservedAtUnixSeconds { 153 /// Validates a positive instant representable by SQLite's signed integer. 154 pub fn new(value: u64) -> Result<Self, RhiTradeMutationAdmissionError> { 155 if value == 0 || i64::try_from(value).is_err() { 156 return Err(failure( 157 RhiTradeMutationAdmissionErrorKind::InvalidObservationTime, 158 )); 159 } 160 Ok(Self(value)) 161 } 162 163 /// Returns the injected observation time. 164 #[must_use] 165 pub const fn get(self) -> u64 { 166 self.0 167 } 168 } 169 170 /// Explicit caller-selected future authored-time tolerance with no default. 171 #[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)] 172 pub struct RhiTradeMutationAuthoredTimePolicy { 173 maximum_future_seconds: u64, 174 } 175 176 impl RhiTradeMutationAuthoredTimePolicy { 177 /// Validates the inclusive maximum future skew. 178 pub fn new(maximum_future_seconds: u64) -> Result<Self, RhiTradeMutationAdmissionError> { 179 if i64::try_from(maximum_future_seconds).is_err() { 180 return Err(failure( 181 RhiTradeMutationAdmissionErrorKind::InvalidTimePolicy, 182 )); 183 } 184 Ok(Self { 185 maximum_future_seconds, 186 }) 187 } 188 189 /// Returns the inclusive maximum future skew. 190 #[must_use] 191 pub const fn maximum_future_seconds(self) -> u64 { 192 self.maximum_future_seconds 193 } 194 } 195 196 /// Stable source-free classification for trade-mutation admission failures. 197 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 198 pub enum RhiTradeMutationAdmissionErrorKind { 199 InvalidLimits, 200 EmptyEvent, 201 EventTooLarge, 202 InvalidEventUtf8, 203 MalformedEvent, 204 DuplicateEventField, 205 EventIdentifierTooLarge, 206 EventContentTooLarge, 207 TooManyTags, 208 TooManyTagElements, 209 TagElementTooLarge, 210 TagsTooLarge, 211 TooManyExtraFields, 212 ExtraFieldsTooLarge, 213 InvalidObservationTime, 214 InvalidTimePolicy, 215 InvalidAuthoredTime, 216 InvalidEventId, 217 InvalidSignature, 218 UnsupportedKind, 219 InvalidAuthor, 220 InvalidMutation, 221 AuthoredTimeRejected, 222 } 223 224 impl RhiTradeMutationAdmissionErrorKind { 225 const fn message(self) -> &'static str { 226 match self { 227 Self::InvalidLimits => "trade-event admission limits are invalid", 228 Self::EmptyEvent => "trade event bytes are empty", 229 Self::EventTooLarge => "trade event exceeds its wire limit", 230 Self::InvalidEventUtf8 => "trade event is not valid UTF-8", 231 Self::MalformedEvent => "trade event structure is invalid", 232 Self::DuplicateEventField => "trade event contains a duplicate field", 233 Self::EventIdentifierTooLarge => "trade event identifier exceeds its limit", 234 Self::EventContentTooLarge => "trade event content exceeds its limit", 235 Self::TooManyTags => "trade event tag count exceeds its limit", 236 Self::TooManyTagElements => "trade event tag elements exceed their count limit", 237 Self::TagElementTooLarge => "trade event tag element exceeds its byte limit", 238 Self::TagsTooLarge => "trade event tags exceed their aggregate byte limit", 239 Self::TooManyExtraFields => "trade event extras exceed their field limit", 240 Self::ExtraFieldsTooLarge => "trade event extras exceed their byte limit", 241 Self::InvalidObservationTime => "trade-event observation time is invalid", 242 Self::InvalidTimePolicy => "trade-event authored-time policy is invalid", 243 Self::InvalidAuthoredTime => "trade-event authored time is invalid", 244 Self::InvalidEventId => "trade event identifier verification failed", 245 Self::InvalidSignature => "trade event signature verification failed", 246 Self::UnsupportedKind => "trade event kind is unsupported", 247 Self::InvalidAuthor => "trade event author binding is invalid", 248 Self::InvalidMutation => "trade mutation contract is invalid", 249 Self::AuthoredTimeRejected => "trade-event authored time is outside policy", 250 } 251 } 252 } 253 254 /// One redacted trade-mutation admission failure. 255 #[derive(Clone, Copy, PartialEq, Eq)] 256 pub struct RhiTradeMutationAdmissionError { 257 kind: RhiTradeMutationAdmissionErrorKind, 258 } 259 260 impl RhiTradeMutationAdmissionError { 261 /// Returns the stable failure classification. 262 #[must_use] 263 pub const fn kind(self) -> RhiTradeMutationAdmissionErrorKind { 264 self.kind 265 } 266 } 267 268 impl fmt::Debug for RhiTradeMutationAdmissionError { 269 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 270 formatter 271 .debug_struct("RhiTradeMutationAdmissionError") 272 .field("kind", &self.kind) 273 .finish() 274 } 275 } 276 277 impl fmt::Display for RhiTradeMutationAdmissionError { 278 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 279 formatter.write_str(self.kind.message()) 280 } 281 } 282 283 impl Error for RhiTradeMutationAdmissionError {} 284 285 /// One bounded, signature-verified, canonical trade-mutation event. 286 /// 287 /// Construction is sealed to the admission boundary: 288 /// 289 /// ```compile_fail 290 /// use rhi::RhiAdmittedTradeMutationEvent; 291 /// 292 /// let _forged = RhiAdmittedTradeMutationEvent {}; 293 /// ``` 294 pub struct RhiAdmittedTradeMutationEvent { 295 original: Box<[u8]>, 296 event: EventEnvelope, 297 mutation: TradeMutationEnvelopeV1, 298 mutation_id: MutationId, 299 observed_at: RhiTradeMutationObservedAtUnixSeconds, 300 } 301 302 impl RhiAdmittedTradeMutationEvent { 303 /// Returns the exact bounded wire bytes supplied to the admission boundary. 304 #[must_use] 305 pub fn original_bytes(&self) -> &[u8] { 306 &self.original 307 } 308 309 /// Returns the independently verified Nostr event identifier. 310 #[must_use] 311 pub fn event_id(&self) -> &EventId { 312 self.event.id() 313 } 314 315 /// Returns the canonical content-derived mutation identifier. 316 #[must_use] 317 pub const fn mutation_id(&self) -> &MutationId { 318 &self.mutation_id 319 } 320 321 /// Returns the exact registered Nostr event kind. 322 #[must_use] 323 pub fn event_kind(&self) -> u32 { 324 self.event.kind_u32() 325 } 326 327 /// Returns the untrusted-but-policy-admitted event-authored UTC second. 328 #[must_use] 329 pub fn authored_at_unix_seconds(&self) -> u64 { 330 self.event.created_at_u64() 331 } 332 333 /// Returns the injected UTC second used to admit this source observation. 334 #[must_use] 335 pub const fn observed_at_unix_seconds(&self) -> RhiTradeMutationObservedAtUnixSeconds { 336 self.observed_at 337 } 338 339 /// Returns the canonical typed mutation bound to the signed event. 340 #[must_use] 341 pub const fn mutation(&self) -> &TradeMutationEnvelopeV1 { 342 &self.mutation 343 } 344 345 pub(crate) fn event_signature_bytes(&self) -> [u8; 64] { 346 *self.event.sig().as_bytes() 347 } 348 349 pub(crate) fn into_parts( 350 self, 351 ) -> ( 352 Box<[u8]>, 353 EventEnvelope, 354 TradeMutationEnvelopeV1, 355 MutationId, 356 RhiTradeMutationObservedAtUnixSeconds, 357 ) { 358 ( 359 self.original, 360 self.event, 361 self.mutation, 362 self.mutation_id, 363 self.observed_at, 364 ) 365 } 366 } 367 368 impl fmt::Debug for RhiAdmittedTradeMutationEvent { 369 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 370 formatter 371 .debug_struct("RhiAdmittedTradeMutationEvent") 372 .field("wire_bytes", &self.original.len()) 373 .field("content_bytes", &self.event.content().len()) 374 .field("tag_count", &self.event.tags().len()) 375 .field("event_kind", &self.event.kind_u32()) 376 .field("authored_at_unix_seconds", &self.event.created_at_u64()) 377 .field("identity", &"[redacted]") 378 .finish() 379 } 380 } 381 382 /// Bounds, verifies, and admits one canonical signed trade-mutation event. 383 pub fn admit_rhi_trade_mutation_event( 384 limits: RhiTradeMutationAdmissionLimits, 385 original: &[u8], 386 observed_at: RhiTradeMutationObservedAtUnixSeconds, 387 authored_time_policy: RhiTradeMutationAuthoredTimePolicy, 388 ) -> Result<RhiAdmittedTradeMutationEvent, RhiTradeMutationAdmissionError> { 389 if original.is_empty() { 390 return Err(failure(RhiTradeMutationAdmissionErrorKind::EmptyEvent)); 391 } 392 if original.len() > limits.wire_bytes { 393 return Err(failure(RhiTradeMutationAdmissionErrorKind::EventTooLarge)); 394 } 395 let source = std::str::from_utf8(original) 396 .map_err(|_| failure(RhiTradeMutationAdmissionErrorKind::InvalidEventUtf8))?; 397 preflight_wire(source, limits)?; 398 399 let wire = Nip01EventWire::parse_json_unverified_with_limits(source, limits.wire_limits()) 400 .map_err(|_| failure(RhiTradeMutationAdmissionErrorKind::MalformedEvent))?; 401 let event = wire 402 .into_unverified_envelope() 403 .map_err(|_| failure(RhiTradeMutationAdmissionErrorKind::MalformedEvent))?; 404 405 match verify_id(&event) { 406 Verification::IdVerified => {} 407 Verification::IdMismatch => { 408 return Err(failure(RhiTradeMutationAdmissionErrorKind::InvalidEventId)); 409 } 410 _ => return Err(failure(RhiTradeMutationAdmissionErrorKind::MalformedEvent)), 411 } 412 match verify(&event) { 413 Verification::Verified => {} 414 Verification::IdMismatch => { 415 return Err(failure(RhiTradeMutationAdmissionErrorKind::InvalidEventId)); 416 } 417 Verification::SignatureInvalid => { 418 return Err(failure( 419 RhiTradeMutationAdmissionErrorKind::InvalidSignature, 420 )); 421 } 422 Verification::IdVerified | Verification::MalformedEnvelope => { 423 return Err(failure(RhiTradeMutationAdmissionErrorKind::MalformedEvent)); 424 } 425 } 426 427 validate_authored_time_representation(event.created_at_u64())?; 428 if !is_trade_mutation_event_kind(event.kind_u32()) { 429 return Err(failure(RhiTradeMutationAdmissionErrorKind::UnsupportedKind)); 430 } 431 let mutation = trade_mutation_from_event(&event).map_err(classify_mutation_error)?; 432 let mutation_id = mutation 433 .mutation_id 434 .ok_or_else(|| failure(RhiTradeMutationAdmissionErrorKind::InvalidMutation))?; 435 enforce_authored_time_policy(event.created_at_u64(), observed_at, authored_time_policy)?; 436 437 Ok(RhiAdmittedTradeMutationEvent { 438 original: original.into(), 439 event, 440 mutation, 441 mutation_id, 442 observed_at, 443 }) 444 } 445 446 fn validate_authored_time_representation( 447 authored_at: u64, 448 ) -> Result<(), RhiTradeMutationAdmissionError> { 449 if i64::try_from(authored_at).is_err() { 450 return Err(failure( 451 RhiTradeMutationAdmissionErrorKind::InvalidAuthoredTime, 452 )); 453 } 454 Ok(()) 455 } 456 457 fn enforce_authored_time_policy( 458 authored_at: u64, 459 observed_at: RhiTradeMutationObservedAtUnixSeconds, 460 policy: RhiTradeMutationAuthoredTimePolicy, 461 ) -> Result<(), RhiTradeMutationAdmissionError> { 462 let latest = observed_at 463 .get() 464 .saturating_add(policy.maximum_future_seconds()); 465 if authored_at > latest { 466 return Err(failure( 467 RhiTradeMutationAdmissionErrorKind::AuthoredTimeRejected, 468 )); 469 } 470 Ok(()) 471 } 472 473 fn classify_mutation_error(error: RadrootsTradeMutationError) -> RhiTradeMutationAdmissionError { 474 let kind = match error { 475 RadrootsTradeMutationError::InvalidKind => { 476 RhiTradeMutationAdmissionErrorKind::UnsupportedKind 477 } 478 RadrootsTradeMutationError::AuthorMismatch => { 479 RhiTradeMutationAdmissionErrorKind::InvalidAuthor 480 } 481 _ => RhiTradeMutationAdmissionErrorKind::InvalidMutation, 482 }; 483 failure(kind) 484 } 485 486 fn config_limit( 487 configuration: &RhiConfigDocumentV1, 488 pointer: &str, 489 ) -> Result<usize, RhiTradeMutationAdmissionError> { 490 configuration 491 .normalized() 492 .pointer(pointer) 493 .and_then(serde_json::Value::as_u64) 494 .and_then(|value| usize::try_from(value).ok()) 495 .filter(|value| *value > 0) 496 .ok_or_else(|| failure(RhiTradeMutationAdmissionErrorKind::InvalidLimits)) 497 } 498 499 struct RawEvent<'a> { 500 id: &'a RawValue, 501 pubkey: &'a RawValue, 502 created_at: &'a RawValue, 503 kind: &'a RawValue, 504 tags: &'a RawValue, 505 content: &'a RawValue, 506 sig: &'a RawValue, 507 } 508 509 fn preflight_wire( 510 source: &str, 511 limits: RhiTradeMutationAdmissionLimits, 512 ) -> Result<(), RhiTradeMutationAdmissionError> { 513 let mut deserializer = serde_json::Deserializer::from_str(source); 514 let raw = RawEventSeed 515 .deserialize(&mut deserializer) 516 .map_err(classify_wire_preflight_error)?; 517 deserializer 518 .end() 519 .map_err(|_| failure(RhiTradeMutationAdmissionErrorKind::MalformedEvent))?; 520 521 validate_bounded_string( 522 raw.id, 523 RHI_TRADE_EVENT_ID_MAX_BYTES, 524 RhiTradeMutationAdmissionErrorKind::EventIdentifierTooLarge, 525 )?; 526 validate_bounded_string( 527 raw.pubkey, 528 RHI_TRADE_EVENT_PUBLIC_KEY_MAX_BYTES, 529 RhiTradeMutationAdmissionErrorKind::EventIdentifierTooLarge, 530 )?; 531 validate_bounded_string( 532 raw.sig, 533 RHI_TRADE_EVENT_SIGNATURE_MAX_BYTES, 534 RhiTradeMutationAdmissionErrorKind::EventIdentifierTooLarge, 535 )?; 536 validate_bounded_string( 537 raw.content, 538 limits.content_bytes, 539 RhiTradeMutationAdmissionErrorKind::EventContentTooLarge, 540 )?; 541 parse_scalar::<u64>(raw.created_at)?; 542 parse_scalar::<u32>(raw.kind)?; 543 measure_tags(raw.tags, limits)?; 544 Ok(()) 545 } 546 547 struct RawEventSeed; 548 549 impl<'de> DeserializeSeed<'de> for RawEventSeed { 550 type Value = RawEvent<'de>; 551 552 fn deserialize<D>(self, deserializer: D) -> Result<Self::Value, D::Error> 553 where 554 D: serde::Deserializer<'de>, 555 { 556 deserializer.deserialize_map(RawEventVisitor) 557 } 558 } 559 560 struct RawEventVisitor; 561 562 impl<'de> Visitor<'de> for RawEventVisitor { 563 type Value = RawEvent<'de>; 564 565 fn expecting(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 566 formatter.write_str("a bounded NIP-01 event object") 567 } 568 569 fn visit_map<A>(self, mut map: A) -> Result<Self::Value, A::Error> 570 where 571 A: MapAccess<'de>, 572 { 573 let mut id = None; 574 let mut pubkey = None; 575 let mut created_at = None; 576 let mut kind = None; 577 let mut tags = None; 578 let mut content = None; 579 let mut sig = None; 580 let mut extras = BTreeSet::new(); 581 let mut extra_bytes = 0usize; 582 583 while let Some(key) = map.next_key::<Cow<'de, str>>()? { 584 let slot = match key.as_ref() { 585 "id" => Some(&mut id), 586 "pubkey" => Some(&mut pubkey), 587 "created_at" => Some(&mut created_at), 588 "kind" => Some(&mut kind), 589 "tags" => Some(&mut tags), 590 "content" => Some(&mut content), 591 "sig" => Some(&mut sig), 592 _ => None, 593 }; 594 if let Some(slot) = slot { 595 if slot.is_some() { 596 return Err(de::Error::custom(DUPLICATE_FIELD_SENTINEL)); 597 } 598 *slot = Some(map.next_value::<&'de RawValue>()?); 599 continue; 600 } 601 602 if !extras.insert(key.clone()) { 603 return Err(de::Error::custom(DUPLICATE_FIELD_SENTINEL)); 604 } 605 if extras.len() > RHI_TRADE_EVENT_EXTRA_FIELD_MAX_COUNT { 606 return Err(de::Error::custom(EXTRA_COUNT_SENTINEL)); 607 } 608 let value = map.next_value::<&'de RawValue>()?; 609 extra_bytes = extra_bytes 610 .checked_add(canonical_json_string_len(key.as_ref())) 611 .and_then(|total| total.checked_add(1)) 612 .and_then(|total| total.checked_add(value.get().len())) 613 .ok_or_else(|| de::Error::custom(EXTRA_BYTES_SENTINEL))?; 614 if extra_bytes > RHI_TRADE_EVENT_EXTRA_JSON_MAX_BYTES { 615 return Err(de::Error::custom(EXTRA_BYTES_SENTINEL)); 616 } 617 } 618 619 Ok(RawEvent { 620 id: required(id)?, 621 pubkey: required(pubkey)?, 622 created_at: required(created_at)?, 623 kind: required(kind)?, 624 tags: required(tags)?, 625 content: required(content)?, 626 sig: required(sig)?, 627 }) 628 } 629 } 630 631 fn required<E>(value: Option<&RawValue>) -> Result<&RawValue, E> 632 where 633 E: de::Error, 634 { 635 value.ok_or_else(|| de::Error::custom("missing required event field")) 636 } 637 638 fn canonical_json_string_len(value: &str) -> usize { 639 value.chars().fold(2usize, |length, character| { 640 length.saturating_add(match character { 641 '"' | '\\' | '\n' | '\r' | '\t' | '\u{08}' | '\u{0c}' => 2, 642 '\u{00}'..='\u{1f}' => 6, 643 _ => character.len_utf8(), 644 }) 645 }) 646 } 647 648 fn parse_scalar<T>(raw: &RawValue) -> Result<T, RhiTradeMutationAdmissionError> 649 where 650 T: serde::de::DeserializeOwned, 651 { 652 serde_json::from_str(raw.get()) 653 .map_err(|_| failure(RhiTradeMutationAdmissionErrorKind::MalformedEvent)) 654 } 655 656 fn validate_bounded_string( 657 raw: &RawValue, 658 maximum: usize, 659 too_large: RhiTradeMutationAdmissionErrorKind, 660 ) -> Result<usize, RhiTradeMutationAdmissionError> { 661 let length = decoded_json_string_utf8_bytes(raw.get()) 662 .ok_or_else(|| failure(RhiTradeMutationAdmissionErrorKind::MalformedEvent))?; 663 if length > maximum { 664 return Err(failure(too_large)); 665 } 666 Ok(length) 667 } 668 669 fn decoded_json_string_utf8_bytes(raw: &str) -> Option<usize> { 670 let bytes = raw.as_bytes(); 671 if bytes.len() < 2 || bytes.first() != Some(&b'"') || bytes.last() != Some(&b'"') { 672 return None; 673 } 674 let end = bytes.len() - 1; 675 let mut index = 1; 676 let mut length = 0usize; 677 while index < end { 678 let byte = bytes[index]; 679 if byte == b'\\' { 680 index = index.checked_add(1)?; 681 let escaped = *bytes.get(index)?; 682 match escaped { 683 b'"' | b'\\' | b'/' | b'b' | b'f' | b'n' | b'r' | b't' => { 684 length = length.checked_add(1)?; 685 index = index.checked_add(1)?; 686 } 687 b'u' => { 688 let first = parse_hex_u16(bytes.get(index + 1..index + 5)?)?; 689 index = index.checked_add(5)?; 690 let scalar = if (0xd800..=0xdbff).contains(&first) { 691 if bytes.get(index..index + 2)? != b"\\u" { 692 return None; 693 } 694 let second = parse_hex_u16(bytes.get(index + 2..index + 6)?)?; 695 if !(0xdc00..=0xdfff).contains(&second) { 696 return None; 697 } 698 index = index.checked_add(6)?; 699 0x1_0000 700 + ((u32::from(first) - 0xd800) << 10) 701 + (u32::from(second) - 0xdc00) 702 } else if (0xdc00..=0xdfff).contains(&first) { 703 return None; 704 } else { 705 u32::from(first) 706 }; 707 length = length.checked_add(char::from_u32(scalar)?.len_utf8())?; 708 } 709 _ => return None, 710 } 711 } else if byte < 0x80 { 712 if byte < 0x20 || byte == b'"' { 713 return None; 714 } 715 length = length.checked_add(1)?; 716 index = index.checked_add(1)?; 717 } else { 718 let character = raw.get(index..end)?.chars().next()?; 719 let width = character.len_utf8(); 720 length = length.checked_add(width)?; 721 index = index.checked_add(width)?; 722 } 723 } 724 (index == end).then_some(length) 725 } 726 727 fn parse_hex_u16(bytes: &[u8]) -> Option<u16> { 728 if bytes.len() != 4 { 729 return None; 730 } 731 bytes.iter().try_fold(0u16, |value, byte| { 732 let digit = match byte { 733 b'0'..=b'9' => u16::from(byte - b'0'), 734 b'a'..=b'f' => u16::from(byte - b'a') + 10, 735 b'A'..=b'F' => u16::from(byte - b'A') + 10, 736 _ => return None, 737 }; 738 value.checked_mul(16)?.checked_add(digit) 739 }) 740 } 741 742 #[derive(Clone, Copy)] 743 struct Measurement { 744 count: usize, 745 elements: usize, 746 bytes: usize, 747 } 748 749 fn measure_tags( 750 raw: &RawValue, 751 limits: RhiTradeMutationAdmissionLimits, 752 ) -> Result<Measurement, RhiTradeMutationAdmissionError> { 753 let mut deserializer = serde_json::Deserializer::from_str(raw.get()); 754 let measurement = TagsSeed { limits } 755 .deserialize(&mut deserializer) 756 .map_err(classify_tag_error)?; 757 deserializer 758 .end() 759 .map_err(|_| failure(RhiTradeMutationAdmissionErrorKind::MalformedEvent))?; 760 Ok(measurement) 761 } 762 763 struct TagsSeed { 764 limits: RhiTradeMutationAdmissionLimits, 765 } 766 767 impl<'de> DeserializeSeed<'de> for TagsSeed { 768 type Value = Measurement; 769 770 fn deserialize<D>(self, deserializer: D) -> Result<Self::Value, D::Error> 771 where 772 D: serde::Deserializer<'de>, 773 { 774 deserializer.deserialize_seq(TagsVisitor { 775 limits: self.limits, 776 }) 777 } 778 } 779 780 struct TagsVisitor { 781 limits: RhiTradeMutationAdmissionLimits, 782 } 783 784 impl<'de> Visitor<'de> for TagsVisitor { 785 type Value = Measurement; 786 787 fn expecting(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 788 formatter.write_str("a bounded array of Nostr tags") 789 } 790 791 fn visit_seq<A>(self, mut sequence: A) -> Result<Self::Value, A::Error> 792 where 793 A: SeqAccess<'de>, 794 { 795 let mut result = Measurement { 796 count: 0, 797 elements: 0, 798 bytes: 0, 799 }; 800 while result.count < self.limits.tag_count { 801 let remaining_elements = self 802 .limits 803 .tag_total_elements 804 .checked_sub(result.elements) 805 .ok_or_else(|| de::Error::custom(TAG_ELEMENT_COUNT_SENTINEL))?; 806 let Some(tag) = sequence.next_element_seed(TagSeed { 807 maximum_elements: remaining_elements, 808 maximum_element_bytes: self.limits.tag_element_bytes, 809 })? 810 else { 811 return Ok(result); 812 }; 813 result.count += 1; 814 result.elements = result 815 .elements 816 .checked_add(tag.elements) 817 .ok_or_else(|| de::Error::custom(TAG_ELEMENT_COUNT_SENTINEL))?; 818 result.bytes = result 819 .bytes 820 .checked_add(tag.bytes) 821 .ok_or_else(|| de::Error::custom(TAG_TOTAL_BYTES_SENTINEL))?; 822 if result.bytes > self.limits.tag_total_bytes { 823 return Err(de::Error::custom(TAG_TOTAL_BYTES_SENTINEL)); 824 } 825 } 826 if sequence.next_element::<IgnoredAny>()?.is_some() { 827 return Err(de::Error::custom(TAG_COUNT_SENTINEL)); 828 } 829 Ok(result) 830 } 831 } 832 833 struct TagSeed { 834 maximum_elements: usize, 835 maximum_element_bytes: usize, 836 } 837 838 impl<'de> DeserializeSeed<'de> for TagSeed { 839 type Value = Measurement; 840 841 fn deserialize<D>(self, deserializer: D) -> Result<Self::Value, D::Error> 842 where 843 D: serde::Deserializer<'de>, 844 { 845 deserializer.deserialize_seq(TagVisitor { 846 maximum_elements: self.maximum_elements, 847 maximum_element_bytes: self.maximum_element_bytes, 848 }) 849 } 850 } 851 852 struct TagVisitor { 853 maximum_elements: usize, 854 maximum_element_bytes: usize, 855 } 856 857 impl<'de> Visitor<'de> for TagVisitor { 858 type Value = Measurement; 859 860 fn expecting(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 861 formatter.write_str("a bounded Nostr tag") 862 } 863 864 fn visit_seq<A>(self, mut sequence: A) -> Result<Self::Value, A::Error> 865 where 866 A: SeqAccess<'de>, 867 { 868 let mut result = Measurement { 869 count: 1, 870 elements: 0, 871 bytes: 0, 872 }; 873 while result.elements < self.maximum_elements { 874 let Some(length) = sequence.next_element_seed(StringLengthSeed { 875 maximum: self.maximum_element_bytes, 876 })? 877 else { 878 return Ok(result); 879 }; 880 result.elements += 1; 881 result.bytes = result 882 .bytes 883 .checked_add(length) 884 .ok_or_else(|| de::Error::custom(TAG_TOTAL_BYTES_SENTINEL))?; 885 } 886 if sequence.next_element::<IgnoredAny>()?.is_some() { 887 return Err(de::Error::custom(TAG_ELEMENT_COUNT_SENTINEL)); 888 } 889 Ok(result) 890 } 891 } 892 893 struct StringLengthSeed { 894 maximum: usize, 895 } 896 897 impl<'de> DeserializeSeed<'de> for StringLengthSeed { 898 type Value = usize; 899 900 fn deserialize<D>(self, deserializer: D) -> Result<Self::Value, D::Error> 901 where 902 D: serde::Deserializer<'de>, 903 { 904 let raw = <&RawValue>::deserialize(deserializer)?; 905 let length = decoded_json_string_utf8_bytes(raw.get()) 906 .ok_or_else(|| de::Error::invalid_type(de::Unexpected::Other("non-string"), &self))?; 907 if length > self.maximum { 908 return Err(de::Error::custom(TAG_ELEMENT_BYTES_SENTINEL)); 909 } 910 Ok(length) 911 } 912 } 913 914 impl de::Expected for StringLengthSeed { 915 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 916 formatter.write_str("a JSON string") 917 } 918 } 919 920 fn classify_wire_preflight_error(error: serde_json::Error) -> RhiTradeMutationAdmissionError { 921 let rendered = error.to_string(); 922 let kind = if rendered.contains(DUPLICATE_FIELD_SENTINEL) { 923 RhiTradeMutationAdmissionErrorKind::DuplicateEventField 924 } else if rendered.contains(EXTRA_COUNT_SENTINEL) { 925 RhiTradeMutationAdmissionErrorKind::TooManyExtraFields 926 } else if rendered.contains(EXTRA_BYTES_SENTINEL) { 927 RhiTradeMutationAdmissionErrorKind::ExtraFieldsTooLarge 928 } else { 929 RhiTradeMutationAdmissionErrorKind::MalformedEvent 930 }; 931 failure(kind) 932 } 933 934 fn classify_tag_error(error: serde_json::Error) -> RhiTradeMutationAdmissionError { 935 let rendered = error.to_string(); 936 let kind = if rendered.contains(TAG_COUNT_SENTINEL) { 937 RhiTradeMutationAdmissionErrorKind::TooManyTags 938 } else if rendered.contains(TAG_ELEMENT_COUNT_SENTINEL) { 939 RhiTradeMutationAdmissionErrorKind::TooManyTagElements 940 } else if rendered.contains(TAG_ELEMENT_BYTES_SENTINEL) { 941 RhiTradeMutationAdmissionErrorKind::TagElementTooLarge 942 } else if rendered.contains(TAG_TOTAL_BYTES_SENTINEL) { 943 RhiTradeMutationAdmissionErrorKind::TagsTooLarge 944 } else { 945 RhiTradeMutationAdmissionErrorKind::MalformedEvent 946 }; 947 failure(kind) 948 } 949 950 const fn failure(kind: RhiTradeMutationAdmissionErrorKind) -> RhiTradeMutationAdmissionError { 951 RhiTradeMutationAdmissionError { kind } 952 } 953 954 #[cfg(test)] 955 mod tests { 956 use super::*; 957 958 #[test] 959 fn canonical_json_string_length_is_allocation_free_and_exact() { 960 assert_eq!(canonical_json_string_len("plain"), 7); 961 assert_eq!(canonical_json_string_len("a\nb"), 6); 962 assert_eq!(canonical_json_string_len("é"), 4); 963 assert_eq!(canonical_json_string_len("\u{0001}"), 8); 964 } 965 966 #[test] 967 fn decoded_json_string_length_handles_escapes_and_surrogates() { 968 assert_eq!(decoded_json_string_utf8_bytes(r#""plain""#), Some(5)); 969 assert_eq!(decoded_json_string_utf8_bytes(r#""a\nb""#), Some(3)); 970 assert_eq!(decoded_json_string_utf8_bytes(r#""\u00e9""#), Some(2)); 971 assert_eq!(decoded_json_string_utf8_bytes(r#""\ud83c\udf31""#), Some(4)); 972 assert_eq!(decoded_json_string_utf8_bytes(r#""\ud83c""#), None); 973 } 974 }