lib

Core libraries for Radroots
git clone https://radroots.dev/git/lib.git
Log | Files | Refs | README

source.rs (47879B)


      1 //! Inbound event source SPI and bounded page models.
      2 
      3 use crate::{
      4     Error, TransportId,
      5     outcome::FetchTargetOutcome,
      6     target::{TargetFingerprint, TargetSet},
      7 };
      8 use alloc::{
      9     boxed::Box,
     10     collections::{BTreeMap, BTreeSet},
     11     string::String,
     12     vec::Vec,
     13 };
     14 use core::{fmt, future::Future, pin::Pin};
     15 use radroots_event::SignedEvent;
     16 use radroots_identity::PublicKey;
     17 
     18 #[cfg(feature = "serde")]
     19 use crate::target::TARGET_SET_MAX_ITEMS;
     20 
     21 pub use crate::status::SourceStatus;
     22 
     23 /// Maximum encoded request identity length.
     24 pub const FETCH_REQUEST_ID_MAX_BYTES: usize = 256;
     25 /// Maximum opaque cursor length.
     26 pub const FETCH_CURSOR_MAX_BYTES: usize = 2_048;
     27 /// Maximum number of events one page may request.
     28 pub const FETCH_PAGE_MAX_EVENTS: u16 = 1_000;
     29 /// Maximum distinct event kinds in one source selector.
     30 pub const FETCH_SELECTOR_MAX_KINDS: usize = 64;
     31 /// Maximum distinct event authors in one source selector.
     32 pub const FETCH_SELECTOR_MAX_AUTHORS: usize = 256;
     33 /// Maximum distinct exact single-letter tag keys in one source selector.
     34 pub const FETCH_SELECTOR_MAX_TAG_KEYS: usize = 26;
     35 /// Maximum exact tag values across one source selector.
     36 pub const FETCH_SELECTOR_MAX_TAG_VALUES: usize = 256;
     37 /// Maximum UTF-8 bytes in one exact tag value.
     38 pub const FETCH_SELECTOR_TAG_VALUE_MAX_BYTES: usize = 4_096;
     39 
     40 // Deliberate representation indirection keeps every selector-bearing request
     41 // and terminal value compact while allocating nothing for the common no-tag
     42 // case. The map itself still owns its bounded tree nodes.
     43 #[allow(clippy::box_collection)]
     44 type ExactTagFilters = Box<BTreeMap<char, Vec<String>>>;
     45 
     46 /// Maximum encoded live-subscription request identity length.
     47 pub const SUBSCRIPTION_REQUEST_ID_MAX_BYTES: usize = 256;
     48 /// Maximum number of events one live subscription may emit.
     49 pub const SUBSCRIPTION_MAX_EVENTS: u16 = 1_000;
     50 
     51 /// Heap-backed future returned by transport SPIs.
     52 pub type BoxFuture<'a, T> = Pin<Box<dyn Future<Output = T> + Send + 'a>>;
     53 
     54 /// Validated caller identity for one fetch operation.
     55 #[derive(Clone, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
     56 pub struct FetchRequestId(String);
     57 
     58 impl FetchRequestId {
     59     /// Parses a non-empty, bounded, printable request identity.
     60     pub fn parse(value: impl Into<String>) -> Result<Self, Error> {
     61         let value = value.into();
     62         if value.is_empty() {
     63             return Err(Error::EmptyFetchRequestId);
     64         }
     65         if value.len() > FETCH_REQUEST_ID_MAX_BYTES
     66             || value != value.trim()
     67             || value.chars().any(char::is_control)
     68         {
     69             return Err(Error::InvalidFetchRequestId);
     70         }
     71         Ok(Self(value))
     72     }
     73 
     74     /// Returns the validated request identity.
     75     pub fn as_str(&self) -> &str {
     76         self.0.as_str()
     77     }
     78 }
     79 
     80 impl fmt::Display for FetchRequestId {
     81     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
     82         formatter.write_str(self.as_str())
     83     }
     84 }
     85 
     86 /// Opaque adapter-owned continuation token.
     87 #[derive(Clone, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
     88 pub struct FetchCursor(String);
     89 
     90 impl FetchCursor {
     91     /// Parses a bounded printable cursor without interpreting its contents.
     92     pub fn parse(value: impl Into<String>) -> Result<Self, Error> {
     93         let value = value.into();
     94         if value.is_empty() {
     95             return Err(Error::EmptyFetchCursor);
     96         }
     97         if value.len() > FETCH_CURSOR_MAX_BYTES
     98             || value != value.trim()
     99             || value.chars().any(char::is_control)
    100         {
    101             return Err(Error::InvalidFetchCursor);
    102         }
    103         Ok(Self(value))
    104     }
    105 
    106     /// Returns the opaque cursor exactly as supplied by its adapter.
    107     pub fn as_str(&self) -> &str {
    108         self.0.as_str()
    109     }
    110 }
    111 
    112 impl fmt::Display for FetchCursor {
    113     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    114         formatter.write_str(self.as_str())
    115     }
    116 }
    117 
    118 /// Validated caller identity for one bounded live subscription.
    119 #[derive(Clone, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
    120 pub struct SubscriptionRequestId(String);
    121 
    122 impl SubscriptionRequestId {
    123     /// Parses a non-empty, bounded, printable request identity.
    124     pub fn parse(value: impl AsRef<str>) -> Result<Self, Error> {
    125         let value = value.as_ref();
    126         if value.is_empty() {
    127             return Err(Error::EmptySubscriptionRequestId);
    128         }
    129         if value.len() > SUBSCRIPTION_REQUEST_ID_MAX_BYTES
    130             || value != value.trim()
    131             || value.chars().any(char::is_control)
    132         {
    133             return Err(Error::InvalidSubscriptionRequestId);
    134         }
    135         Ok(Self(String::from(value)))
    136     }
    137 
    138     /// Returns the validated request identity.
    139     pub fn as_str(&self) -> &str {
    140         self.0.as_str()
    141     }
    142 }
    143 
    144 impl fmt::Display for SubscriptionRequestId {
    145     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    146         formatter.write_str(self.as_str())
    147     }
    148 }
    149 
    150 /// Hard bounds for one source operation.
    151 #[cfg_attr(feature = "serde", derive(serde::Serialize))]
    152 #[derive(Clone, Copy, Debug, Eq, PartialEq)]
    153 pub struct FetchBounds {
    154     limit: u16,
    155     deadline_unix_ms: u64,
    156 }
    157 
    158 impl FetchBounds {
    159     /// Creates bounds with a non-zero page limit and absolute deadline.
    160     pub const fn new(limit: u16, deadline_unix_ms: u64) -> Result<Self, Error> {
    161         if limit == 0 || limit > FETCH_PAGE_MAX_EVENTS {
    162             return Err(Error::InvalidFetchLimit);
    163         }
    164         if deadline_unix_ms == 0 {
    165             return Err(Error::InvalidFetchDeadline);
    166         }
    167         Ok(Self {
    168             limit,
    169             deadline_unix_ms,
    170         })
    171     }
    172 
    173     /// Maximum number of events the adapter may return.
    174     pub const fn limit(self) -> u16 {
    175         self.limit
    176     }
    177 
    178     /// Absolute Unix deadline in milliseconds.
    179     pub const fn deadline_unix_ms(self) -> u64 {
    180         self.deadline_unix_ms
    181     }
    182 }
    183 
    184 /// Hard bounds for one live-subscription operation.
    185 #[cfg_attr(feature = "serde", derive(serde::Serialize))]
    186 #[derive(Clone, Copy, Debug, Eq, PartialEq)]
    187 pub struct SubscriptionBounds {
    188     event_limit: u16,
    189     deadline_unix_ms: u64,
    190 }
    191 
    192 impl SubscriptionBounds {
    193     /// Creates bounds with a non-zero event limit and absolute deadline.
    194     pub const fn new(event_limit: u16, deadline_unix_ms: u64) -> Result<Self, Error> {
    195         if event_limit == 0 || event_limit > SUBSCRIPTION_MAX_EVENTS {
    196             return Err(Error::InvalidSubscriptionLimit);
    197         }
    198         if deadline_unix_ms == 0 {
    199             return Err(Error::InvalidSubscriptionDeadline);
    200         }
    201         Ok(Self {
    202             event_limit,
    203             deadline_unix_ms,
    204         })
    205     }
    206 
    207     /// Maximum number of events the adapter may emit.
    208     pub const fn event_limit(self) -> u16 {
    209         self.event_limit
    210     }
    211 
    212     /// Absolute Unix deadline in milliseconds.
    213     pub const fn deadline_unix_ms(self) -> u64 {
    214         self.deadline_unix_ms
    215     }
    216 }
    217 
    218 /// Transport-neutral constraints applied before a source page is bounded.
    219 ///
    220 /// An empty kind, author, or tag collection means "any" for that dimension.
    221 /// Values for one tag key are alternatives, while distinct tag keys are
    222 /// conjunctive. Time bounds are inclusive Unix seconds. Adapters must apply
    223 /// every configured dimension remotely when their protocol supports it and
    224 /// must defensively exclude non-matching events before returning a page.
    225 #[cfg_attr(feature = "serde", derive(serde::Serialize))]
    226 #[derive(Clone, Debug, Default, Eq, PartialEq)]
    227 pub struct FetchSelector {
    228     kinds: Vec<u32>,
    229     authors: Vec<PublicKey>,
    230     #[cfg_attr(
    231         feature = "serde",
    232         serde(serialize_with = "serde_impl::serialize_exact_tags")
    233     )]
    234     exact_tags: Option<ExactTagFilters>,
    235     since_unix_seconds: Option<u64>,
    236     until_unix_seconds: Option<u64>,
    237 }
    238 
    239 impl FetchSelector {
    240     /// Creates a selector that accepts every event within request bounds.
    241     #[must_use]
    242     pub const fn all() -> Self {
    243         Self {
    244             kinds: Vec::new(),
    245             authors: Vec::new(),
    246             exact_tags: None,
    247             since_unix_seconds: None,
    248             until_unix_seconds: None,
    249         }
    250     }
    251 
    252     /// Restricts the selector to exact, unique event kinds.
    253     pub fn with_kinds(mut self, mut kinds: Vec<u32>) -> Result<Self, Error> {
    254         if kinds.len() > FETCH_SELECTOR_MAX_KINDS {
    255             return Err(Error::FetchSelectorTooLarge);
    256         }
    257         kinds.sort_unstable();
    258         if kinds.windows(2).any(|pair| pair[0] == pair[1]) {
    259             return Err(Error::DuplicateFetchKind);
    260         }
    261         self.kinds = kinds;
    262         Ok(self)
    263     }
    264 
    265     /// Restricts the selector to exact, unique canonical authors.
    266     pub fn with_authors(mut self, mut authors: Vec<PublicKey>) -> Result<Self, Error> {
    267         if authors.len() > FETCH_SELECTOR_MAX_AUTHORS {
    268             return Err(Error::FetchSelectorTooLarge);
    269         }
    270         authors.sort();
    271         if authors.windows(2).any(|pair| pair[0] == pair[1]) {
    272             return Err(Error::DuplicateFetchAuthor);
    273         }
    274         self.authors = authors;
    275         Ok(self)
    276     }
    277 
    278     /// Requires one exact indexed single-letter tag value.
    279     ///
    280     /// Repeating a key adds an alternative value for that key. Different keys
    281     /// are conjunctive. Keys are lowercase ASCII letters and values are
    282     /// non-empty bounded UTF-8 strings.
    283     pub fn with_exact_tag_value(
    284         mut self,
    285         key: char,
    286         value: impl AsRef<str>,
    287     ) -> Result<Self, Error> {
    288         if !key.is_ascii_lowercase() {
    289             return Err(Error::InvalidFetchTagKey);
    290         }
    291         let value = value.as_ref();
    292         if value.is_empty()
    293             || value.len() > FETCH_SELECTOR_TAG_VALUE_MAX_BYTES
    294             || value.chars().any(char::is_control)
    295         {
    296             return Err(Error::InvalidFetchTagValue);
    297         }
    298         let exact_tags = self.exact_tags.as_deref();
    299         let total_values = exact_tags
    300             .into_iter()
    301             .flat_map(BTreeMap::values)
    302             .map(Vec::len)
    303             .sum::<usize>();
    304         if (!exact_tags.is_some_and(|tags| tags.contains_key(&key))
    305             && exact_tags.is_some_and(|tags| tags.len() == FETCH_SELECTOR_MAX_TAG_KEYS))
    306             || total_values == FETCH_SELECTOR_MAX_TAG_VALUES
    307         {
    308             return Err(Error::FetchSelectorTooLarge);
    309         }
    310         let values = self
    311             .exact_tags
    312             .get_or_insert_with(|| Box::new(BTreeMap::new()))
    313             .entry(key)
    314             .or_default();
    315         match values.binary_search_by(|candidate| candidate.as_str().cmp(value)) {
    316             Ok(_) => return Err(Error::DuplicateFetchTagValue),
    317             Err(position) => values.insert(position, String::from(value)),
    318         }
    319         Ok(self)
    320     }
    321 
    322     /// Sets an inclusive lower event-time bound.
    323     pub fn with_since_unix_seconds(mut self, since: u64) -> Result<Self, Error> {
    324         if self.until_unix_seconds.is_some_and(|until| since > until) {
    325             return Err(Error::InvalidFetchTimeRange);
    326         }
    327         self.since_unix_seconds = Some(since);
    328         Ok(self)
    329     }
    330 
    331     /// Sets an inclusive upper event-time bound.
    332     pub fn with_until_unix_seconds(mut self, until: u64) -> Result<Self, Error> {
    333         if self.since_unix_seconds.is_some_and(|since| since > until) {
    334             return Err(Error::InvalidFetchTimeRange);
    335         }
    336         self.until_unix_seconds = Some(until);
    337         Ok(self)
    338     }
    339 
    340     /// Returns sorted exact event kinds, or an empty slice for any kind.
    341     pub fn kinds(&self) -> &[u32] {
    342         self.kinds.as_slice()
    343     }
    344 
    345     /// Returns sorted exact authors, or an empty slice for any author.
    346     pub fn authors(&self) -> &[PublicKey] {
    347         self.authors.as_slice()
    348     }
    349 
    350     /// Returns exact tag filters in canonical key order.
    351     pub fn exact_tag_filters(&self) -> impl Iterator<Item = (char, &[String])> + '_ {
    352         self.exact_tags
    353             .iter()
    354             .flat_map(|tags| tags.iter())
    355             .map(|(key, values)| (*key, values.as_slice()))
    356     }
    357 
    358     /// Returns the inclusive lower event-time bound.
    359     pub const fn since_unix_seconds(&self) -> Option<u64> {
    360         self.since_unix_seconds
    361     }
    362 
    363     /// Returns the inclusive upper event-time bound.
    364     pub const fn until_unix_seconds(&self) -> Option<u64> {
    365         self.until_unix_seconds
    366     }
    367 
    368     /// Returns whether one canonical signed event satisfies every dimension.
    369     #[must_use]
    370     pub fn matches(&self, event: &SignedEvent) -> bool {
    371         (self.kinds.is_empty() || self.kinds.binary_search(&event.kind()).is_ok())
    372             && (self.authors.is_empty() || self.authors.binary_search(event.pubkey()).is_ok())
    373             && self.exact_tags.as_deref().is_none_or(|tags| {
    374                 tags.iter().all(|(key, values)| {
    375                     event.envelope().tag_slices().iter().any(|tag| {
    376                         let elements = tag.as_slice();
    377                         elements.first().is_some_and(|candidate| {
    378                             candidate.len() == 1 && candidate.starts_with(*key)
    379                         }) && elements.get(1).is_some_and(|value| {
    380                             values
    381                                 .binary_search_by(|candidate| candidate.as_str().cmp(value))
    382                                 .is_ok()
    383                         })
    384                     })
    385                 })
    386             })
    387             && self
    388                 .since_unix_seconds
    389                 .is_none_or(|since| event.created_at() >= since)
    390             && self
    391                 .until_unix_seconds
    392                 .is_none_or(|until| event.created_at() <= until)
    393     }
    394 }
    395 
    396 /// Bounded request for one page from one or more transport targets.
    397 #[cfg_attr(feature = "serde", derive(serde::Serialize))]
    398 #[derive(Clone, Debug, Eq, PartialEq)]
    399 pub struct FetchRequest {
    400     request_id: FetchRequestId,
    401     target_set: TargetSet,
    402     bounds: FetchBounds,
    403     cursor: Option<FetchCursor>,
    404     selector: FetchSelector,
    405 }
    406 
    407 impl FetchRequest {
    408     /// Creates a first-page request.
    409     pub fn new(
    410         request_id: impl Into<String>,
    411         target_set: TargetSet,
    412         bounds: FetchBounds,
    413     ) -> Result<Self, Error> {
    414         Ok(Self {
    415             request_id: FetchRequestId::parse(request_id)?,
    416             target_set,
    417             bounds,
    418             cursor: None,
    419             selector: FetchSelector::all(),
    420         })
    421     }
    422 
    423     /// Sets the adapter-owned cursor for a continuation request.
    424     #[must_use]
    425     pub fn with_cursor(mut self, cursor: FetchCursor) -> Self {
    426         self.cursor = Some(cursor);
    427         self
    428     }
    429 
    430     /// Applies explicit transport-neutral event constraints.
    431     #[must_use]
    432     pub fn with_selector(mut self, selector: FetchSelector) -> Self {
    433         self.selector = selector;
    434         self
    435     }
    436 
    437     /// Returns the request identity.
    438     pub fn request_id(&self) -> &FetchRequestId {
    439         &self.request_id
    440     }
    441 
    442     /// Returns the exact requested target set.
    443     pub fn target_set(&self) -> &TargetSet {
    444         &self.target_set
    445     }
    446 
    447     /// Returns the hard operation bounds.
    448     pub const fn bounds(&self) -> FetchBounds {
    449         self.bounds
    450     }
    451 
    452     /// Returns the continuation cursor, when this is not a first-page request.
    453     pub fn cursor(&self) -> Option<&FetchCursor> {
    454         self.cursor.as_ref()
    455     }
    456 
    457     /// Returns the exact event constraints for this request.
    458     pub const fn selector(&self) -> &FetchSelector {
    459         &self.selector
    460     }
    461 }
    462 
    463 /// Transport observation attached to one inbound event.
    464 #[cfg_attr(feature = "serde", derive(serde::Serialize))]
    465 #[cfg_attr(feature = "serde", serde(deny_unknown_fields))]
    466 #[derive(Clone, Debug, Eq, PartialEq)]
    467 pub struct EventProvenance {
    468     transport_id: TransportId,
    469     target: TargetFingerprint,
    470     observed_at_unix_ms: u64,
    471     cursor: Option<FetchCursor>,
    472 }
    473 
    474 impl EventProvenance {
    475     /// Creates provenance for an event observed from one exact target.
    476     pub fn new(
    477         transport_id: TransportId,
    478         target: TargetFingerprint,
    479         observed_at_unix_ms: u64,
    480     ) -> Result<Self, Error> {
    481         if observed_at_unix_ms == 0 {
    482             return Err(Error::InvalidObservedAt);
    483         }
    484         Ok(Self {
    485             transport_id,
    486             target,
    487             observed_at_unix_ms,
    488             cursor: None,
    489         })
    490     }
    491 
    492     /// Attaches the adapter cursor that located this event.
    493     #[must_use]
    494     pub fn with_cursor(mut self, cursor: FetchCursor) -> Self {
    495         self.cursor = Some(cursor);
    496         self
    497     }
    498 
    499     /// Returns the transport that produced this observation.
    500     pub const fn transport_id(&self) -> TransportId {
    501         self.transport_id
    502     }
    503 
    504     /// Returns the exact target fingerprint that produced this observation.
    505     pub const fn target(&self) -> &TargetFingerprint {
    506         &self.target
    507     }
    508 
    509     /// Returns the host-recorded observation time.
    510     pub const fn observed_at_unix_ms(&self) -> u64 {
    511         self.observed_at_unix_ms
    512     }
    513 
    514     /// Returns the optional adapter cursor at the observation point.
    515     pub const fn cursor(&self) -> Option<&FetchCursor> {
    516         self.cursor.as_ref()
    517     }
    518 }
    519 
    520 /// ID-checked signed event plus transport provenance.
    521 ///
    522 /// Signature verification, contract validation, canonical admission, storage,
    523 /// and projection results intentionally remain outside this transport model.
    524 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
    525 #[cfg_attr(feature = "serde", serde(deny_unknown_fields))]
    526 #[derive(Clone, Debug, Eq, PartialEq)]
    527 pub struct ObservedEvent {
    528     event: SignedEvent,
    529     provenance: EventProvenance,
    530 }
    531 
    532 impl ObservedEvent {
    533     /// Attaches provenance to an inbound signed event.
    534     pub const fn new(event: SignedEvent, provenance: EventProvenance) -> Self {
    535         Self { event, provenance }
    536     }
    537 
    538     /// Returns the unverified signed event payload.
    539     pub const fn event(&self) -> &SignedEvent {
    540         &self.event
    541     }
    542 
    543     /// Returns the transport observation.
    544     pub const fn provenance(&self) -> &EventProvenance {
    545         &self.provenance
    546     }
    547 }
    548 
    549 /// State required to continue or conclude a bounded fetch.
    550 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
    551 #[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))]
    552 #[derive(Clone, Debug, Eq, PartialEq)]
    553 pub enum NextPage {
    554     /// Every requested target reached its current end.
    555     Complete,
    556     /// More results are available from this exact cursor.
    557     Cursor(FetchCursor),
    558     /// The operation was cancelled and may optionally be resumed.
    559     Cancelled { resume_from: Option<FetchCursor> },
    560 }
    561 
    562 /// One validated, request-bound page of inbound observations.
    563 #[cfg_attr(feature = "serde", derive(serde::Serialize))]
    564 #[derive(Clone, Debug, Eq, PartialEq)]
    565 pub struct FetchPage {
    566     request_id: FetchRequestId,
    567     target_set: TargetSet,
    568     limit: u16,
    569     selector: FetchSelector,
    570     events: Vec<ObservedEvent>,
    571     target_outcomes: Vec<FetchTargetOutcome>,
    572     next_page: NextPage,
    573 }
    574 
    575 impl FetchPage {
    576     /// Creates and validates one page against its originating request.
    577     pub fn for_request(
    578         request: &FetchRequest,
    579         events: Vec<ObservedEvent>,
    580         target_outcomes: Vec<FetchTargetOutcome>,
    581         next_page: NextPage,
    582     ) -> Result<Self, Error> {
    583         let page = Self {
    584             request_id: request.request_id.clone(),
    585             target_set: request.target_set.clone(),
    586             limit: request.bounds.limit,
    587             selector: request.selector.clone(),
    588             events,
    589             target_outcomes,
    590             next_page,
    591         };
    592         page.validate()?;
    593         Ok(page)
    594     }
    595 
    596     /// Validates internal cardinality, target provenance, and outcome identity.
    597     pub fn validate(&self) -> Result<(), Error> {
    598         if self.limit == 0 || self.limit > FETCH_PAGE_MAX_EVENTS {
    599             return Err(Error::InvalidFetchLimit);
    600         }
    601         if self.events.len() > usize::from(self.limit) {
    602             return Err(Error::FetchPageLimitExceeded);
    603         }
    604 
    605         for observed in &self.events {
    606             if !self.selector.matches(observed.event()) {
    607                 return Err(Error::UnexpectedFetchEvent);
    608             }
    609             let provenance = observed.provenance();
    610             let Some(target) = self
    611                 .target_set
    612                 .targets()
    613                 .iter()
    614                 .find(|target| target.fingerprint() == provenance.target())
    615             else {
    616                 return Err(Error::UnexpectedFetchProvenance);
    617             };
    618             if *target.kind() != provenance.transport_id() {
    619                 return Err(Error::UnexpectedFetchProvenance);
    620             }
    621         }
    622 
    623         let requested: BTreeSet<&str> = self
    624             .target_set
    625             .targets()
    626             .iter()
    627             .map(|target| target.fingerprint().as_str())
    628             .collect();
    629         let mut outcomes = BTreeSet::new();
    630         for outcome in &self.target_outcomes {
    631             if !requested.contains(outcome.target().as_str()) {
    632                 return Err(Error::UnexpectedFetchTargetOutcome);
    633             }
    634             if !outcomes.insert(outcome.target().as_str()) {
    635                 return Err(Error::DuplicateFetchTargetOutcome);
    636             }
    637         }
    638         Ok(())
    639     }
    640 
    641     /// Validates that this page is bound to the exact originating request.
    642     pub fn validate_for_request(&self, request: &FetchRequest) -> Result<(), Error> {
    643         self.validate()?;
    644         if &self.request_id != request.request_id()
    645             || self.target_set != *request.target_set()
    646             || self.limit != request.bounds().limit()
    647             || self.selector != *request.selector()
    648         {
    649             return Err(Error::FetchPageRequestMismatch);
    650         }
    651         Ok(())
    652     }
    653 
    654     /// Returns the request identity.
    655     pub const fn request_id(&self) -> &FetchRequestId {
    656         &self.request_id
    657     }
    658 
    659     /// Returns the observations in adapter order.
    660     pub fn events(&self) -> &[ObservedEvent] {
    661         self.events.as_slice()
    662     }
    663 
    664     /// Returns zero or more target-specific outcomes; omitted targets remain unreported.
    665     pub fn target_outcomes(&self) -> &[FetchTargetOutcome] {
    666         self.target_outcomes.as_slice()
    667     }
    668 
    669     /// Returns continuation, completion, or cancellation state.
    670     pub const fn next_page(&self) -> &NextPage {
    671         &self.next_page
    672     }
    673 }
    674 
    675 /// Per-target continuation point for a live subscription.
    676 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
    677 #[cfg_attr(feature = "serde", serde(deny_unknown_fields))]
    678 #[derive(Clone, Debug, Eq, PartialEq)]
    679 pub struct SubscriptionCheckpoint {
    680     target: TargetFingerprint,
    681     cursor: FetchCursor,
    682 }
    683 
    684 impl SubscriptionCheckpoint {
    685     /// Binds an opaque adapter cursor to one exact target.
    686     pub const fn new(target: TargetFingerprint, cursor: FetchCursor) -> Self {
    687         Self { target, cursor }
    688     }
    689 
    690     /// Returns the exact target fingerprint.
    691     pub const fn target(&self) -> &TargetFingerprint {
    692         &self.target
    693     }
    694 
    695     /// Returns the opaque adapter cursor.
    696     pub const fn cursor(&self) -> &FetchCursor {
    697         &self.cursor
    698     }
    699 }
    700 
    701 /// Bounded request for a live stream from one or more transport targets.
    702 #[cfg_attr(feature = "serde", derive(serde::Serialize))]
    703 #[derive(Clone, Debug, Eq, PartialEq)]
    704 pub struct SubscriptionRequest {
    705     request_id: SubscriptionRequestId,
    706     target_set: TargetSet,
    707     bounds: SubscriptionBounds,
    708     selector: FetchSelector,
    709     checkpoints: Vec<SubscriptionCheckpoint>,
    710 }
    711 
    712 impl SubscriptionRequest {
    713     /// Creates a live request with no prior target checkpoints.
    714     pub fn new(
    715         request_id: impl AsRef<str>,
    716         target_set: TargetSet,
    717         bounds: SubscriptionBounds,
    718     ) -> Result<Self, Error> {
    719         Ok(Self {
    720             request_id: SubscriptionRequestId::parse(request_id)?,
    721             target_set,
    722             bounds,
    723             selector: FetchSelector::all(),
    724             checkpoints: Vec::new(),
    725         })
    726     }
    727 
    728     /// Applies explicit transport-neutral event constraints.
    729     #[must_use]
    730     pub fn with_selector(mut self, selector: FetchSelector) -> Self {
    731         self.selector = selector;
    732         self
    733     }
    734 
    735     /// Applies a bounded, unique checkpoint subset in canonical target order.
    736     pub fn with_checkpoints<I>(mut self, checkpoints: I) -> Result<Self, Error>
    737     where
    738         I: IntoIterator<Item = SubscriptionCheckpoint>,
    739     {
    740         self.checkpoints = normalize_subscription_checkpoints(&self.target_set, checkpoints)?;
    741         Ok(self)
    742     }
    743 
    744     /// Returns the request identity.
    745     pub const fn request_id(&self) -> &SubscriptionRequestId {
    746         &self.request_id
    747     }
    748 
    749     /// Returns the exact requested target set.
    750     pub const fn target_set(&self) -> &TargetSet {
    751         &self.target_set
    752     }
    753 
    754     /// Returns the hard operation bounds.
    755     pub const fn bounds(&self) -> SubscriptionBounds {
    756         self.bounds
    757     }
    758 
    759     /// Returns the exact event constraints for this request.
    760     pub const fn selector(&self) -> &FetchSelector {
    761         &self.selector
    762     }
    763 
    764     /// Returns prior per-target checkpoints in canonical target order.
    765     pub fn checkpoints(&self) -> &[SubscriptionCheckpoint] {
    766         self.checkpoints.as_slice()
    767     }
    768 }
    769 
    770 /// One request-bound live event and its resulting per-target checkpoint.
    771 #[cfg_attr(feature = "serde", derive(serde::Serialize))]
    772 #[derive(Clone, Debug, Eq, PartialEq)]
    773 pub struct SubscriptionEvent {
    774     request: SubscriptionRequest,
    775     observed: ObservedEvent,
    776     checkpoint: SubscriptionCheckpoint,
    777 }
    778 
    779 impl SubscriptionEvent {
    780     /// Creates and validates one live event against its originating request.
    781     pub fn for_request(
    782         request: &SubscriptionRequest,
    783         observed: ObservedEvent,
    784         checkpoint: SubscriptionCheckpoint,
    785     ) -> Result<Self, Error> {
    786         let event = Self {
    787             request: request.clone(),
    788             observed,
    789             checkpoint,
    790         };
    791         event.validate_for_request(request)?;
    792         Ok(event)
    793     }
    794 
    795     /// Validates target, transport, selector, cursor, and request identity.
    796     pub fn validate_for_request(&self, request: &SubscriptionRequest) -> Result<(), Error> {
    797         if self.request != *request {
    798             return Err(Error::UnexpectedSubscriptionEvent);
    799         }
    800         let provenance = self.observed.provenance();
    801         let Some(target) = request
    802             .target_set
    803             .targets()
    804             .iter()
    805             .find(|target| target.fingerprint() == provenance.target())
    806         else {
    807             return Err(Error::UnexpectedSubscriptionEvent);
    808         };
    809         if *target.kind() != provenance.transport_id()
    810             || !request.selector.matches(self.observed.event())
    811         {
    812             return Err(Error::UnexpectedSubscriptionEvent);
    813         }
    814         if self.checkpoint.target() != provenance.target()
    815             || provenance.cursor() != Some(self.checkpoint.cursor())
    816         {
    817             return Err(Error::SubscriptionEventCheckpointMismatch);
    818         }
    819         Ok(())
    820     }
    821 
    822     /// Returns the request identity.
    823     pub const fn request_id(&self) -> &SubscriptionRequestId {
    824         self.request.request_id()
    825     }
    826 
    827     /// Returns the exact originating request.
    828     pub const fn request(&self) -> &SubscriptionRequest {
    829         &self.request
    830     }
    831 
    832     /// Returns the observed event.
    833     pub const fn observed(&self) -> &ObservedEvent {
    834         &self.observed
    835     }
    836 
    837     /// Returns the checkpoint established by this event.
    838     pub const fn checkpoint(&self) -> &SubscriptionCheckpoint {
    839         &self.checkpoint
    840     }
    841 }
    842 
    843 /// Stable terminal reason for a bounded live subscription.
    844 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
    845 #[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))]
    846 #[derive(Clone, Copy, Debug, Eq, PartialEq)]
    847 pub enum SubscriptionEndReason {
    848     /// The requested maximum event count was emitted.
    849     EventLimit,
    850     /// The absolute request deadline was reached.
    851     Deadline,
    852     /// Explicit or future-drop cancellation was observed.
    853     Cancelled,
    854     /// The underlying source closed before another event was available.
    855     SourceClosed,
    856 }
    857 
    858 /// Request-bound terminal result for a live subscription.
    859 #[cfg_attr(feature = "serde", derive(serde::Serialize))]
    860 #[derive(Clone, Debug, Eq, PartialEq)]
    861 pub struct SubscriptionEnd {
    862     request: SubscriptionRequest,
    863     event_count: u16,
    864     checkpoints: Vec<SubscriptionCheckpoint>,
    865     reason: SubscriptionEndReason,
    866 }
    867 
    868 impl SubscriptionEnd {
    869     /// Creates a terminal result with canonical final checkpoints.
    870     pub fn for_request<I>(
    871         request: &SubscriptionRequest,
    872         event_count: u16,
    873         checkpoints: I,
    874         reason: SubscriptionEndReason,
    875     ) -> Result<Self, Error>
    876     where
    877         I: IntoIterator<Item = SubscriptionCheckpoint>,
    878     {
    879         if event_count > request.bounds.event_limit {
    880             return Err(Error::SubscriptionEndLimitExceeded);
    881         }
    882         if reason == SubscriptionEndReason::EventLimit && event_count != request.bounds.event_limit
    883         {
    884             return Err(Error::InvalidSubscriptionEnd);
    885         }
    886         Ok(Self {
    887             request: request.clone(),
    888             event_count,
    889             checkpoints: normalize_subscription_checkpoints(&request.target_set, checkpoints)?,
    890             reason,
    891         })
    892     }
    893 
    894     /// Validates that this result belongs to the exact originating request.
    895     pub fn validate_for_request(&self, request: &SubscriptionRequest) -> Result<(), Error> {
    896         if self.request != *request || self.event_count > request.bounds.event_limit {
    897             return Err(Error::SubscriptionEndRequestMismatch);
    898         }
    899         normalize_subscription_checkpoints(&request.target_set, self.checkpoints.clone())?;
    900         Ok(())
    901     }
    902 
    903     /// Returns the exact originating request.
    904     pub const fn request(&self) -> &SubscriptionRequest {
    905         &self.request
    906     }
    907 
    908     /// Returns the number of events emitted before termination.
    909     pub const fn event_count(&self) -> u16 {
    910         self.event_count
    911     }
    912 
    913     /// Returns final per-target checkpoints in canonical target order.
    914     pub fn checkpoints(&self) -> &[SubscriptionCheckpoint] {
    915         self.checkpoints.as_slice()
    916     }
    917 
    918     /// Returns why the bounded operation terminated.
    919     pub const fn reason(&self) -> SubscriptionEndReason {
    920         self.reason
    921     }
    922 }
    923 
    924 /// One event or the stable terminal result from a live subscription.
    925 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
    926 #[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))]
    927 #[derive(Clone, Debug, Eq, PartialEq)]
    928 pub enum SubscriptionNext {
    929     /// One request-bound event was observed.
    930     Event(Box<SubscriptionEvent>),
    931     /// The subscription reached its stable terminal state.
    932     End(SubscriptionEnd),
    933 }
    934 
    935 fn normalize_subscription_checkpoints<I>(
    936     target_set: &TargetSet,
    937     checkpoints: I,
    938 ) -> Result<Vec<SubscriptionCheckpoint>, Error>
    939 where
    940     I: IntoIterator<Item = SubscriptionCheckpoint>,
    941 {
    942     let maximum = target_set.len();
    943     let mut checkpoints: Vec<_> = checkpoints.into_iter().take(maximum + 1).collect();
    944     if checkpoints.len() > maximum {
    945         return Err(Error::SubscriptionCheckpointSetTooLarge);
    946     }
    947 
    948     let mut seen = BTreeSet::new();
    949     for checkpoint in &checkpoints {
    950         if target_position(target_set, checkpoint.target()).is_none() {
    951             return Err(Error::UnexpectedSubscriptionCheckpoint);
    952         }
    953         if !seen.insert(checkpoint.target().as_str()) {
    954             return Err(Error::DuplicateSubscriptionCheckpoint);
    955         }
    956     }
    957     checkpoints.sort_by_key(|checkpoint| {
    958         target_position(target_set, checkpoint.target()).expect("checkpoint target validated")
    959     });
    960     Ok(checkpoints)
    961 }
    962 
    963 fn target_position(target_set: &TargetSet, fingerprint: &TargetFingerprint) -> Option<usize> {
    964     target_set
    965         .targets()
    966         .iter()
    967         .position(|target| target.fingerprint() == fingerprint)
    968 }
    969 
    970 /// Host SPI for inbound event retrieval.
    971 ///
    972 /// This trait supports external implementations and is dyn-compatible. Its
    973 /// futures are `Send`; implementations must not borrow request data after a
    974 /// future completes. `status` observes source state and does not initiate a
    975 /// fetch. `fetch` performs at most the work bounded by its request and owns no
    976 /// hidden retry loop.
    977 ///
    978 /// Dropping a returned future requests cancellation. If it is dropped before
    979 /// a remote request is published, the implementation must leave no remote
    980 /// operation behind. Once publication may have occurred, cancellation cannot
    981 /// claim rollback; a later observation may report the remote outcome. The
    982 /// explicit request deadline bounds work independently of future cancellation.
    983 pub trait EventSource: Send + Sync {
    984     /// Returns the source's current runtime status.
    985     fn status(&self) -> BoxFuture<'_, Result<SourceStatus, Error>>;
    986 
    987     /// Fetches one bounded page of transport-neutral events.
    988     fn fetch(&self, request: FetchRequest) -> BoxFuture<'_, Result<FetchPage, Error>>;
    989 }
    990 
    991 /// One active, bounded live-subscription operation.
    992 ///
    993 /// Implementations must enforce the request's event limit and absolute
    994 /// deadline without hidden retries. Once [`SubscriptionNext::End`] has been
    995 /// returned, every later `next` or `cancel` call must return the exact same
    996 /// terminal result. Dropping a pending future requests cancellation but does
    997 /// not claim that an already-observed remote event was rolled back.
    998 pub trait EventSubscription: Send {
    999     /// Returns the exact request governing this operation.
   1000     fn request(&self) -> &SubscriptionRequest;
   1001 
   1002     /// Returns the next request-bound event or stable terminal result.
   1003     fn next(&mut self) -> BoxFuture<'_, Result<SubscriptionNext, Error>>;
   1004 
   1005     /// Requests cancellation and returns the stable terminal result.
   1006     fn cancel(&mut self) -> BoxFuture<'_, Result<SubscriptionEnd, Error>>;
   1007 }
   1008 
   1009 /// Heap-owned live-subscription capability returned by adapters.
   1010 pub type BoxSubscription = Box<dyn EventSubscription>;
   1011 
   1012 /// Host SPI for beginning bounded live subscriptions.
   1013 ///
   1014 /// This is separate from [`EventSource`] so existing bounded-fetch producers
   1015 /// remain source-compatible until they explicitly adopt live delivery.
   1016 pub trait EventSubscriber: Send + Sync {
   1017     /// Begins one exact bounded live-subscription request.
   1018     fn subscribe(
   1019         &self,
   1020         request: SubscriptionRequest,
   1021     ) -> BoxFuture<'_, Result<BoxSubscription, Error>>;
   1022 }
   1023 
   1024 #[cfg(feature = "serde")]
   1025 mod serde_impl {
   1026     use super::*;
   1027 
   1028     pub(super) fn serialize_exact_tags<S>(
   1029         exact_tags: &Option<ExactTagFilters>,
   1030         serializer: S,
   1031     ) -> Result<S::Ok, S::Error>
   1032     where
   1033         S: serde::Serializer,
   1034     {
   1035         match exact_tags {
   1036             Some(exact_tags) => serde::Serialize::serialize(exact_tags, serializer),
   1037             None => serde::Serialize::serialize(&BTreeMap::<char, Vec<String>>::new(), serializer),
   1038         }
   1039     }
   1040 
   1041     impl serde::Serialize for FetchRequestId {
   1042         fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
   1043         where
   1044             S: serde::Serializer,
   1045         {
   1046             serializer.serialize_str(self.as_str())
   1047         }
   1048     }
   1049 
   1050     impl<'de> serde::Deserialize<'de> for FetchRequestId {
   1051         fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
   1052         where
   1053             D: serde::Deserializer<'de>,
   1054         {
   1055             let value = <String as serde::Deserialize>::deserialize(deserializer)?;
   1056             Self::parse(value).map_err(serde::de::Error::custom)
   1057         }
   1058     }
   1059 
   1060     impl serde::Serialize for FetchCursor {
   1061         fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
   1062         where
   1063             S: serde::Serializer,
   1064         {
   1065             serializer.serialize_str(self.as_str())
   1066         }
   1067     }
   1068 
   1069     impl<'de> serde::Deserialize<'de> for FetchCursor {
   1070         fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
   1071         where
   1072             D: serde::Deserializer<'de>,
   1073         {
   1074             let value = <String as serde::Deserialize>::deserialize(deserializer)?;
   1075             Self::parse(value).map_err(serde::de::Error::custom)
   1076         }
   1077     }
   1078 
   1079     impl serde::Serialize for SubscriptionRequestId {
   1080         fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
   1081         where
   1082             S: serde::Serializer,
   1083         {
   1084             serializer.serialize_str(self.as_str())
   1085         }
   1086     }
   1087 
   1088     impl<'de> serde::Deserialize<'de> for SubscriptionRequestId {
   1089         fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
   1090         where
   1091             D: serde::Deserializer<'de>,
   1092         {
   1093             let value = <String as serde::Deserialize>::deserialize(deserializer)?;
   1094             Self::parse(value.as_str()).map_err(serde::de::Error::custom)
   1095         }
   1096     }
   1097 
   1098     #[derive(serde::Deserialize)]
   1099     #[serde(deny_unknown_fields)]
   1100     struct FetchBoundsWire {
   1101         limit: u16,
   1102         deadline_unix_ms: u64,
   1103     }
   1104 
   1105     impl<'de> serde::Deserialize<'de> for FetchBounds {
   1106         fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
   1107         where
   1108             D: serde::Deserializer<'de>,
   1109         {
   1110             let wire = FetchBoundsWire::deserialize(deserializer)?;
   1111             Self::new(wire.limit, wire.deadline_unix_ms).map_err(serde::de::Error::custom)
   1112         }
   1113     }
   1114 
   1115     #[derive(serde::Deserialize)]
   1116     #[serde(deny_unknown_fields)]
   1117     struct SubscriptionBoundsWire {
   1118         event_limit: u16,
   1119         deadline_unix_ms: u64,
   1120     }
   1121 
   1122     impl<'de> serde::Deserialize<'de> for SubscriptionBounds {
   1123         fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
   1124         where
   1125             D: serde::Deserializer<'de>,
   1126         {
   1127             let wire = SubscriptionBoundsWire::deserialize(deserializer)?;
   1128             Self::new(wire.event_limit, wire.deadline_unix_ms).map_err(serde::de::Error::custom)
   1129         }
   1130     }
   1131 
   1132     #[derive(Default)]
   1133     struct ExactTagsWire(Vec<(char, Vec<String>)>);
   1134 
   1135     impl<'de> serde::Deserialize<'de> for ExactTagsWire {
   1136         fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
   1137         where
   1138             D: serde::Deserializer<'de>,
   1139         {
   1140             struct ExactTagsVisitor;
   1141 
   1142             impl<'de> serde::de::Visitor<'de> for ExactTagsVisitor {
   1143                 type Value = ExactTagsWire;
   1144 
   1145                 fn expecting(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
   1146                     formatter.write_str("a map of unique exact single-letter tag filters")
   1147                 }
   1148 
   1149                 fn visit_map<A>(self, mut map: A) -> Result<Self::Value, A::Error>
   1150                 where
   1151                     A: serde::de::MapAccess<'de>,
   1152                 {
   1153                     let mut entries = Vec::<(char, Vec<String>)>::new();
   1154                     while let Some((key, values)) = map.next_entry::<char, Vec<String>>()? {
   1155                         if entries.iter().any(|(candidate, _)| *candidate == key) {
   1156                             return Err(serde::de::Error::custom(
   1157                                 "transport fetch selector contains a duplicate tag key",
   1158                             ));
   1159                         }
   1160                         entries.push((key, values));
   1161                     }
   1162                     Ok(ExactTagsWire(entries))
   1163                 }
   1164             }
   1165 
   1166             deserializer.deserialize_map(ExactTagsVisitor)
   1167         }
   1168     }
   1169 
   1170     #[derive(serde::Deserialize)]
   1171     #[serde(deny_unknown_fields)]
   1172     struct FetchRequestWire {
   1173         request_id: String,
   1174         target_set: TargetSet,
   1175         bounds: FetchBounds,
   1176         cursor: Option<FetchCursor>,
   1177         #[serde(default)]
   1178         selector: FetchSelector,
   1179     }
   1180 
   1181     impl<'de> serde::Deserialize<'de> for FetchRequest {
   1182         fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
   1183         where
   1184             D: serde::Deserializer<'de>,
   1185         {
   1186             let wire = FetchRequestWire::deserialize(deserializer)?;
   1187             Self::new(wire.request_id, wire.target_set, wire.bounds)
   1188                 .map(|request| request.with_selector(wire.selector))
   1189                 .map(|request| match wire.cursor {
   1190                     Some(cursor) => request.with_cursor(cursor),
   1191                     None => request,
   1192                 })
   1193                 .map_err(serde::de::Error::custom)
   1194         }
   1195     }
   1196 
   1197     #[derive(serde::Deserialize)]
   1198     #[serde(deny_unknown_fields)]
   1199     struct SubscriptionRequestWire {
   1200         request_id: String,
   1201         target_set: TargetSet,
   1202         bounds: SubscriptionBounds,
   1203         #[serde(default)]
   1204         selector: FetchSelector,
   1205         #[serde(default)]
   1206         #[serde(deserialize_with = "deserialize_subscription_checkpoints")]
   1207         checkpoints: Vec<SubscriptionCheckpoint>,
   1208     }
   1209 
   1210     impl<'de> serde::Deserialize<'de> for SubscriptionRequest {
   1211         fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
   1212         where
   1213             D: serde::Deserializer<'de>,
   1214         {
   1215             let wire = SubscriptionRequestWire::deserialize(deserializer)?;
   1216             Self::new(wire.request_id.as_str(), wire.target_set, wire.bounds)
   1217                 .map(|request| request.with_selector(wire.selector))
   1218                 .and_then(|request| request.with_checkpoints(wire.checkpoints))
   1219                 .map_err(serde::de::Error::custom)
   1220         }
   1221     }
   1222 
   1223     impl<'de> serde::Deserialize<'de> for FetchSelector {
   1224         fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
   1225         where
   1226             D: serde::Deserializer<'de>,
   1227         {
   1228             #[derive(serde::Deserialize)]
   1229             #[serde(deny_unknown_fields)]
   1230             struct Wire {
   1231                 #[serde(default)]
   1232                 kinds: Vec<u32>,
   1233                 #[serde(default)]
   1234                 authors: Vec<PublicKey>,
   1235                 #[serde(default)]
   1236                 exact_tags: ExactTagsWire,
   1237                 since_unix_seconds: Option<u64>,
   1238                 until_unix_seconds: Option<u64>,
   1239             }
   1240 
   1241             let wire = Wire::deserialize(deserializer)?;
   1242             let selector = FetchSelector::all()
   1243                 .with_kinds(wire.kinds)
   1244                 .and_then(|selector| selector.with_authors(wire.authors))
   1245                 .and_then(|selector| {
   1246                     wire.exact_tags
   1247                         .0
   1248                         .into_iter()
   1249                         .try_fold(selector, |selector, (key, values)| {
   1250                             values.into_iter().try_fold(selector, |selector, value| {
   1251                                 selector.with_exact_tag_value(key, value)
   1252                             })
   1253                         })
   1254                 })
   1255                 .and_then(|selector| match wire.since_unix_seconds {
   1256                     Some(since) => selector.with_since_unix_seconds(since),
   1257                     None => Ok(selector),
   1258                 })
   1259                 .and_then(|selector| match wire.until_unix_seconds {
   1260                     Some(until) => selector.with_until_unix_seconds(until),
   1261                     None => Ok(selector),
   1262                 });
   1263             selector.map_err(serde::de::Error::custom)
   1264         }
   1265     }
   1266 
   1267     #[derive(serde::Deserialize)]
   1268     #[serde(deny_unknown_fields)]
   1269     struct EventProvenanceWire {
   1270         transport_id: TransportId,
   1271         target: TargetFingerprint,
   1272         observed_at_unix_ms: u64,
   1273         cursor: Option<FetchCursor>,
   1274     }
   1275 
   1276     impl<'de> serde::Deserialize<'de> for EventProvenance {
   1277         fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
   1278         where
   1279             D: serde::Deserializer<'de>,
   1280         {
   1281             let wire = EventProvenanceWire::deserialize(deserializer)?;
   1282             Self::new(wire.transport_id, wire.target, wire.observed_at_unix_ms)
   1283                 .map(|provenance| match wire.cursor {
   1284                     Some(cursor) => provenance.with_cursor(cursor),
   1285                     None => provenance,
   1286                 })
   1287                 .map_err(serde::de::Error::custom)
   1288         }
   1289     }
   1290 
   1291     #[derive(serde::Deserialize)]
   1292     #[serde(deny_unknown_fields)]
   1293     struct FetchPageWire {
   1294         request_id: FetchRequestId,
   1295         target_set: TargetSet,
   1296         limit: u16,
   1297         #[serde(default)]
   1298         selector: FetchSelector,
   1299         events: Vec<ObservedEvent>,
   1300         target_outcomes: Vec<FetchTargetOutcome>,
   1301         next_page: NextPage,
   1302     }
   1303 
   1304     impl<'de> serde::Deserialize<'de> for FetchPage {
   1305         fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
   1306         where
   1307             D: serde::Deserializer<'de>,
   1308         {
   1309             let wire = FetchPageWire::deserialize(deserializer)?;
   1310             let page = Self {
   1311                 request_id: wire.request_id,
   1312                 target_set: wire.target_set,
   1313                 limit: wire.limit,
   1314                 selector: wire.selector,
   1315                 events: wire.events,
   1316                 target_outcomes: wire.target_outcomes,
   1317                 next_page: wire.next_page,
   1318             };
   1319             page.validate().map_err(serde::de::Error::custom)?;
   1320             Ok(page)
   1321         }
   1322     }
   1323 
   1324     #[derive(serde::Deserialize)]
   1325     #[serde(deny_unknown_fields)]
   1326     struct SubscriptionEventWire {
   1327         request: SubscriptionRequest,
   1328         observed: ObservedEvent,
   1329         checkpoint: SubscriptionCheckpoint,
   1330     }
   1331 
   1332     impl<'de> serde::Deserialize<'de> for SubscriptionEvent {
   1333         fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
   1334         where
   1335             D: serde::Deserializer<'de>,
   1336         {
   1337             let wire = SubscriptionEventWire::deserialize(deserializer)?;
   1338             Self::for_request(&wire.request, wire.observed, wire.checkpoint)
   1339                 .map_err(serde::de::Error::custom)
   1340         }
   1341     }
   1342 
   1343     #[derive(serde::Deserialize)]
   1344     #[serde(deny_unknown_fields)]
   1345     struct SubscriptionEndWire {
   1346         request: SubscriptionRequest,
   1347         event_count: u16,
   1348         #[serde(deserialize_with = "deserialize_subscription_checkpoints")]
   1349         checkpoints: Vec<SubscriptionCheckpoint>,
   1350         reason: SubscriptionEndReason,
   1351     }
   1352 
   1353     impl<'de> serde::Deserialize<'de> for SubscriptionEnd {
   1354         fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
   1355         where
   1356             D: serde::Deserializer<'de>,
   1357         {
   1358             let wire = SubscriptionEndWire::deserialize(deserializer)?;
   1359             Self::for_request(
   1360                 &wire.request,
   1361                 wire.event_count,
   1362                 wire.checkpoints,
   1363                 wire.reason,
   1364             )
   1365             .map_err(serde::de::Error::custom)
   1366         }
   1367     }
   1368 
   1369     fn deserialize_subscription_checkpoints<'de, D>(
   1370         deserializer: D,
   1371     ) -> Result<Vec<SubscriptionCheckpoint>, D::Error>
   1372     where
   1373         D: serde::Deserializer<'de>,
   1374     {
   1375         struct CheckpointVisitor;
   1376 
   1377         impl<'de> serde::de::Visitor<'de> for CheckpointVisitor {
   1378             type Value = Vec<SubscriptionCheckpoint>;
   1379 
   1380             fn expecting(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
   1381                 formatter.write_str("a bounded transport subscription checkpoint sequence")
   1382             }
   1383 
   1384             fn visit_seq<A>(self, mut sequence: A) -> Result<Self::Value, A::Error>
   1385             where
   1386                 A: serde::de::SeqAccess<'de>,
   1387             {
   1388                 let capacity = sequence.size_hint().unwrap_or(0).min(TARGET_SET_MAX_ITEMS);
   1389                 let mut checkpoints = Vec::with_capacity(capacity);
   1390                 while let Some(checkpoint) = sequence.next_element()? {
   1391                     if checkpoints.len() == TARGET_SET_MAX_ITEMS {
   1392                         return Err(serde::de::Error::custom(
   1393                             Error::SubscriptionCheckpointSetTooLarge,
   1394                         ));
   1395                     }
   1396                     checkpoints.push(checkpoint);
   1397                 }
   1398                 Ok(checkpoints)
   1399             }
   1400         }
   1401 
   1402         deserializer.deserialize_seq(CheckpointVisitor)
   1403     }
   1404 }