rhi

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

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 }