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 }