runtime.rs (277239B)
1 #![forbid(unsafe_code)] 2 3 #[cfg(test)] 4 use crate::relay::outbound::protocol_messages_for_test; 5 use crate::{ 6 client_message::RuntimeClientMessage, 7 config::BaseRelayRuntimeConfig, 8 errors::BaseRelayError, 9 event_bus::{TangleEventBus, TangleEventReceiver}, 10 groups::GroupServiceHandle, 11 logging, 12 ops::{BaseRelayReadinessHandle, BaseRelayReadinessState}, 13 pocket_event_validation::{pocket_event_id, pocket_event_kind, pocket_event_pubkey}, 14 rate_limits::{ 15 TangleQueryRateLimitConfig, TangleRateLimitDecision, TangleRateLimitKey, 16 TangleRateLimitQueryClass, TangleRateLimitRule, TangleRateLimitScope, TangleRateLimiter, 17 }, 18 relay::{ 19 auth::BaseAuthState, 20 core::{ 21 BaseRelay, BaseRelayCountQuery, BaseRelayCountReport, BaseRelayEventWrite, 22 BaseRelayFilterLimitMode, BaseRelayLimits, BaseRelayQueryMetrics, BaseRelayQueryReport, 23 BaseRelayReqQuery, BaseRelayShutdownReport, matched_filter_context, 24 }, 25 filter::{BaseRelayMatchedFilterContext, BaseRelayRequestedKinds}, 26 live::LiveSubscriptionSet, 27 outbound::{RuntimeRelayMessage, protocol_control_messages}, 28 }, 29 }; 30 use serde::{Deserialize, Serialize}; 31 use std::{ 32 collections::BTreeSet, 33 fmt, fs, 34 net::IpAddr, 35 num::NonZeroU32, 36 path::Path, 37 str, 38 sync::{ 39 Arc, 40 atomic::{AtomicU64, AtomicUsize, Ordering}, 41 }, 42 time::Instant, 43 }; 44 use tangle_groups::{ 45 GroupAuthContext, GroupEventClass, GroupId, KIND_GROUP_JOIN_REQUEST, StoreOffset, 46 validate_client_group_event_structure, 47 }; 48 use tangle_protocol::{Kind, PublicKeyHex, RelayMessage, SubscriptionId, UnixTimestamp}; 49 use tangle_store_pocket::{ 50 PocketEvent, PocketFilter, PocketOwnedEvent, PocketOwnedFilter, PocketStoreHandle, PocketTime, 51 }; 52 use tokio::sync::watch; 53 54 pub struct RelayRuntime { 55 config: BaseRelayRuntimeConfig, 56 relay: BaseRelay, 57 readiness: BaseRelayReadinessHandle, 58 limits: TangleRuntimeLimits, 59 event_bus: TangleEventBus, 60 rate_limiter: TangleRateLimiter, 61 metrics: TangleRuntimeMetrics, 62 shutdown: TangleShutdownSignal, 63 hooks: Arc<dyn RelayRuntimeHooks>, 64 } 65 66 #[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] 67 pub struct TangleClientRateLimitContext { 68 peer_ip: Option<IpAddr>, 69 connection_id: Option<u64>, 70 } 71 72 impl TangleClientRateLimitContext { 73 pub fn new(peer_ip: Option<IpAddr>, connection_id: Option<u64>) -> Self { 74 Self { 75 peer_ip, 76 connection_id, 77 } 78 } 79 80 pub fn peer_ip(self) -> Option<IpAddr> { 81 self.peer_ip 82 } 83 84 pub fn connection_id(self) -> Option<u64> { 85 self.connection_id 86 } 87 } 88 89 pub trait RelayRuntimeHooks: Send + Sync { 90 fn admit_event(&self, _context: &RelayEventAdmissionContext) -> EventAdmissionDecision { 91 EventAdmissionDecision::Accept 92 } 93 94 fn event_stored(&self, _context: &RelayEventStoredContext) {} 95 96 fn storage_used_bytes(&self) -> Option<u64> { 97 None 98 } 99 100 fn plan_query(&self, _context: &RelayQueryProjectionContext) -> RelayProjectionQueryPlan { 101 RelayProjectionQueryPlan::default() 102 } 103 104 fn live_projection_candidates( 105 &self, 106 _context: &RelayLiveProjectionContext, 107 ) -> Vec<RelayLiveProjectionCandidate> { 108 Vec::new() 109 } 110 111 fn project_event( 112 &self, 113 _context: &RelayEventProjectionContext, 114 ) -> RelayEventProjectionDecision { 115 RelayEventProjectionDecision::Emit 116 } 117 118 fn sanitize_public_message(&self, message: RelayMessage) -> RelayMessage { 119 message 120 } 121 } 122 123 #[derive(Debug, Default)] 124 pub struct NoopRelayRuntimeHooks; 125 126 impl RelayRuntimeHooks for NoopRelayRuntimeHooks {} 127 128 #[derive(Debug, Clone, PartialEq, Eq)] 129 pub enum EventAdmissionDecision { 130 Accept, 131 Reject { message: String }, 132 } 133 134 impl EventAdmissionDecision { 135 pub fn reject(message: impl Into<String>) -> Self { 136 Self::Reject { 137 message: message.into(), 138 } 139 } 140 } 141 142 #[derive(Debug, Clone, Default, PartialEq, Eq)] 143 pub struct RelayProjectionContext { 144 identifier: Option<String>, 145 } 146 147 impl RelayProjectionContext { 148 pub fn named(identifier: impl Into<String>) -> Result<Self, BaseRelayError> { 149 let identifier = identifier.into(); 150 if identifier.is_empty() { 151 return Err(BaseRelayError::invalid( 152 "relay projection identifier must not be empty", 153 )); 154 } 155 if identifier.chars().any(char::is_control) { 156 return Err(BaseRelayError::invalid( 157 "relay projection identifier must not contain control characters", 158 )); 159 } 160 Ok(Self { 161 identifier: Some(identifier), 162 }) 163 } 164 165 pub fn identifier(&self) -> Option<&str> { 166 self.identifier.as_deref() 167 } 168 } 169 170 #[derive(Debug, Clone, Copy, PartialEq, Eq)] 171 pub enum RelayEventProjectionSource { 172 HistoricalQuery, 173 LiveFanout { store_offset: u64 }, 174 } 175 176 impl RelayEventProjectionSource { 177 pub fn store_offset(self) -> Option<u64> { 178 match self { 179 Self::HistoricalQuery => None, 180 Self::LiveFanout { store_offset } => Some(store_offset), 181 } 182 } 183 } 184 185 #[derive(Debug, Clone, Copy, PartialEq, Eq)] 186 pub enum RelayEventProjectionDecision { 187 Emit, 188 Suppress, 189 ReplaceWithStoredOffset { store_offset: u64 }, 190 } 191 192 impl RelayEventProjectionDecision { 193 pub fn replace_with_stored_offset(store_offset: u64) -> Self { 194 Self::ReplaceWithStoredOffset { store_offset } 195 } 196 } 197 198 #[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] 199 pub struct RelayProjectionQueryPlan { 200 limit: RelayProjectionQueryLimit, 201 } 202 203 impl RelayProjectionQueryPlan { 204 pub fn limit_after_projection(candidate_limit: NonZeroU32) -> Self { 205 Self { 206 limit: RelayProjectionQueryLimit::AfterProjection { candidate_limit }, 207 } 208 } 209 210 pub fn limit(&self) -> RelayProjectionQueryLimit { 211 self.limit 212 } 213 } 214 215 #[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] 216 pub enum RelayProjectionQueryLimit { 217 #[default] 218 BeforeProjection, 219 AfterProjection { 220 candidate_limit: NonZeroU32, 221 }, 222 } 223 224 #[derive(Debug, Clone, PartialEq, Eq)] 225 pub struct RelayQueryProjectionContext { 226 subscription_id: SubscriptionId, 227 projection: RelayProjectionContext, 228 filters: Vec<RelayMatchedFilterContext>, 229 authenticated_pubkeys: Vec<PublicKeyHex>, 230 } 231 232 impl RelayQueryProjectionContext { 233 pub fn new( 234 subscription_id: SubscriptionId, 235 projection: RelayProjectionContext, 236 filters: Vec<RelayMatchedFilterContext>, 237 ) -> Self { 238 Self::new_with_authenticated_pubkeys(subscription_id, projection, filters, Vec::new()) 239 } 240 241 pub fn new_with_authenticated_pubkeys( 242 subscription_id: SubscriptionId, 243 projection: RelayProjectionContext, 244 filters: Vec<RelayMatchedFilterContext>, 245 authenticated_pubkeys: Vec<PublicKeyHex>, 246 ) -> Self { 247 Self { 248 subscription_id, 249 projection, 250 filters, 251 authenticated_pubkeys, 252 } 253 } 254 255 pub fn subscription_id(&self) -> &SubscriptionId { 256 &self.subscription_id 257 } 258 259 pub fn projection(&self) -> &RelayProjectionContext { 260 &self.projection 261 } 262 263 pub fn filters(&self) -> &[RelayMatchedFilterContext] { 264 &self.filters 265 } 266 267 pub fn authenticated_pubkeys(&self) -> &[PublicKeyHex] { 268 &self.authenticated_pubkeys 269 } 270 } 271 272 #[derive(Debug, Clone, PartialEq, Eq)] 273 pub struct RelayMatchedFilterContext { 274 filter_index: usize, 275 requested_kinds: RelayRequestedKinds, 276 } 277 278 impl RelayMatchedFilterContext { 279 fn from_base(context: &BaseRelayMatchedFilterContext) -> Self { 280 let requested_kinds = match context.requested_kinds() { 281 BaseRelayRequestedKinds::Absent => RelayRequestedKinds::Absent, 282 BaseRelayRequestedKinds::Explicit(kinds) => { 283 RelayRequestedKinds::Explicit(kinds.clone()) 284 } 285 }; 286 Self { 287 filter_index: context.filter_index(), 288 requested_kinds, 289 } 290 } 291 292 pub fn filter_index(&self) -> usize { 293 self.filter_index 294 } 295 296 pub fn requested_kinds(&self) -> &RelayRequestedKinds { 297 &self.requested_kinds 298 } 299 } 300 301 #[derive(Debug, Clone, PartialEq, Eq)] 302 pub enum RelayRequestedKinds { 303 Absent, 304 Explicit(BTreeSet<u32>), 305 } 306 307 #[derive(Debug, Clone, PartialEq, Eq)] 308 pub struct RelayLiveProjectionContext { 309 projection: RelayProjectionContext, 310 source_store_offset: u64, 311 event: RelayEventContext, 312 authenticated_pubkeys: Vec<PublicKeyHex>, 313 } 314 315 impl RelayLiveProjectionContext { 316 pub fn new( 317 projection: RelayProjectionContext, 318 source_store_offset: u64, 319 event: RelayEventContext, 320 ) -> Self { 321 Self::new_with_authenticated_pubkeys(projection, source_store_offset, event, Vec::new()) 322 } 323 324 pub fn new_with_authenticated_pubkeys( 325 projection: RelayProjectionContext, 326 source_store_offset: u64, 327 event: RelayEventContext, 328 authenticated_pubkeys: Vec<PublicKeyHex>, 329 ) -> Self { 330 Self { 331 projection, 332 source_store_offset, 333 event, 334 authenticated_pubkeys, 335 } 336 } 337 338 pub fn projection(&self) -> &RelayProjectionContext { 339 &self.projection 340 } 341 342 pub fn source_store_offset(&self) -> u64 { 343 self.source_store_offset 344 } 345 346 pub fn event(&self) -> &RelayEventContext { 347 &self.event 348 } 349 350 pub fn authenticated_pubkeys(&self) -> &[PublicKeyHex] { 351 &self.authenticated_pubkeys 352 } 353 } 354 355 #[derive(Debug, Clone, Copy, PartialEq, Eq)] 356 pub struct RelayLiveProjectionCandidate { 357 store_offset: u64, 358 } 359 360 impl RelayLiveProjectionCandidate { 361 pub fn stored_offset(store_offset: u64) -> Self { 362 Self { store_offset } 363 } 364 365 pub fn store_offset(self) -> u64 { 366 self.store_offset 367 } 368 } 369 370 #[derive(Debug, Clone, PartialEq, Eq)] 371 pub struct RelayEventContext { 372 event_id: String, 373 pubkey: String, 374 created_at: u64, 375 kind: u32, 376 tags: Vec<Vec<String>>, 377 content: String, 378 } 379 380 impl RelayEventContext { 381 pub fn new( 382 event_id: String, 383 pubkey: String, 384 created_at: u64, 385 kind: u32, 386 tags: Vec<Vec<String>>, 387 content: String, 388 ) -> Self { 389 Self { 390 event_id, 391 pubkey, 392 created_at, 393 kind, 394 tags, 395 content, 396 } 397 } 398 399 fn from_pocket_event(event: &PocketEvent) -> Result<Self, BaseRelayError> { 400 let tags = event 401 .tags() 402 .map_err(|error| BaseRelayError::invalid(error.to_string()))? 403 .iter() 404 .map(|tag| { 405 tag.map(|value| { 406 str::from_utf8(value) 407 .map(str::to_owned) 408 .map_err(|error| BaseRelayError::invalid(error.to_string())) 409 }) 410 .collect::<Result<Vec<_>, _>>() 411 }) 412 .collect::<Result<Vec<_>, _>>()?; 413 let content = str::from_utf8(event.content()) 414 .map(str::to_owned) 415 .map_err(|error| BaseRelayError::invalid(error.to_string()))?; 416 Ok(Self { 417 event_id: event.id().to_string(), 418 pubkey: event.pubkey().to_string(), 419 created_at: event.created_at().as_u64(), 420 kind: u32::from(event.kind().as_u16()), 421 tags, 422 content, 423 }) 424 } 425 426 pub fn event_id(&self) -> &str { 427 &self.event_id 428 } 429 430 pub fn pubkey(&self) -> &str { 431 &self.pubkey 432 } 433 434 pub fn created_at(&self) -> u64 { 435 self.created_at 436 } 437 438 pub fn kind(&self) -> u32 { 439 self.kind 440 } 441 442 pub fn tags(&self) -> &[Vec<String>] { 443 &self.tags 444 } 445 446 pub fn content(&self) -> &str { 447 &self.content 448 } 449 450 pub fn has_tag(&self, name: &str, value: &str) -> bool { 451 self.tags.iter().any(|tag| { 452 tag.first().is_some_and(|tag_name| tag_name == name) 453 && tag.iter().skip(1).any(|tag_value| tag_value == value) 454 }) 455 } 456 } 457 458 #[derive(Debug, Clone, PartialEq, Eq)] 459 pub struct RelayEventAdmissionContext { 460 event: RelayEventContext, 461 authenticated_pubkeys: Vec<String>, 462 peer_ip: Option<IpAddr>, 463 connection_id: Option<u64>, 464 now: u64, 465 } 466 467 impl RelayEventAdmissionContext { 468 pub fn new( 469 event: RelayEventContext, 470 authenticated_pubkeys: Vec<String>, 471 peer_ip: Option<IpAddr>, 472 connection_id: Option<u64>, 473 now: u64, 474 ) -> Self { 475 Self { 476 event, 477 authenticated_pubkeys, 478 peer_ip, 479 connection_id, 480 now, 481 } 482 } 483 484 pub fn event(&self) -> &RelayEventContext { 485 &self.event 486 } 487 488 pub fn authenticated_pubkeys(&self) -> &[String] { 489 &self.authenticated_pubkeys 490 } 491 492 pub fn peer_ip(&self) -> Option<IpAddr> { 493 self.peer_ip 494 } 495 496 pub fn connection_id(&self) -> Option<u64> { 497 self.connection_id 498 } 499 500 pub fn now(&self) -> u64 { 501 self.now 502 } 503 } 504 505 #[derive(Debug, Clone, PartialEq, Eq)] 506 pub struct RelayEventStoredContext { 507 event: RelayEventContext, 508 store_offsets: Vec<u64>, 509 } 510 511 struct RelayEventProjectionRequest<'a> { 512 subscription_id: &'a SubscriptionId, 513 projection: &'a RelayProjectionContext, 514 source: RelayEventProjectionSource, 515 event: &'a PocketEvent, 516 auth: &'a BaseAuthState, 517 matched_filters: Vec<(BaseRelayMatchedFilterContext, &'a PocketFilter)>, 518 } 519 520 struct RelayLiveProjectionDelivery<'a> { 521 subscriptions: &'a LiveSubscriptionSet, 522 auth: &'a BaseAuthState, 523 projection: &'a RelayProjectionContext, 524 group_auth: &'a GroupAuthContext, 525 delivered: &'a mut BTreeSet<(SubscriptionId, String)>, 526 messages: &'a mut Vec<RuntimeRelayMessage>, 527 } 528 529 #[derive(Debug, Clone, PartialEq, Eq)] 530 pub struct RelayEventProjectionContext { 531 subscription_id: SubscriptionId, 532 projection: RelayProjectionContext, 533 source: RelayEventProjectionSource, 534 matched_filters: Vec<RelayMatchedFilterContext>, 535 event: RelayEventContext, 536 authenticated_pubkeys: Vec<PublicKeyHex>, 537 } 538 539 impl RelayEventProjectionContext { 540 pub fn new( 541 subscription_id: SubscriptionId, 542 projection: RelayProjectionContext, 543 source: RelayEventProjectionSource, 544 matched_filter: RelayMatchedFilterContext, 545 event: RelayEventContext, 546 ) -> Self { 547 Self::new_with_authenticated_pubkeys( 548 subscription_id, 549 projection, 550 source, 551 matched_filter, 552 event, 553 Vec::new(), 554 ) 555 } 556 557 pub fn new_with_authenticated_pubkeys( 558 subscription_id: SubscriptionId, 559 projection: RelayProjectionContext, 560 source: RelayEventProjectionSource, 561 matched_filter: RelayMatchedFilterContext, 562 event: RelayEventContext, 563 authenticated_pubkeys: Vec<PublicKeyHex>, 564 ) -> Self { 565 Self::new_with_matched_filters_and_authenticated_pubkeys( 566 subscription_id, 567 projection, 568 source, 569 vec![matched_filter], 570 event, 571 authenticated_pubkeys, 572 ) 573 } 574 575 pub fn new_with_matched_filters( 576 subscription_id: SubscriptionId, 577 projection: RelayProjectionContext, 578 source: RelayEventProjectionSource, 579 matched_filters: Vec<RelayMatchedFilterContext>, 580 event: RelayEventContext, 581 ) -> Self { 582 Self::new_with_matched_filters_and_authenticated_pubkeys( 583 subscription_id, 584 projection, 585 source, 586 matched_filters, 587 event, 588 Vec::new(), 589 ) 590 } 591 592 pub fn new_with_matched_filters_and_authenticated_pubkeys( 593 subscription_id: SubscriptionId, 594 projection: RelayProjectionContext, 595 source: RelayEventProjectionSource, 596 matched_filters: Vec<RelayMatchedFilterContext>, 597 event: RelayEventContext, 598 authenticated_pubkeys: Vec<PublicKeyHex>, 599 ) -> Self { 600 Self { 601 subscription_id, 602 projection, 603 source, 604 matched_filters, 605 event, 606 authenticated_pubkeys, 607 } 608 } 609 610 pub fn subscription_id(&self) -> &SubscriptionId { 611 &self.subscription_id 612 } 613 614 pub fn projection(&self) -> &RelayProjectionContext { 615 &self.projection 616 } 617 618 pub fn source(&self) -> RelayEventProjectionSource { 619 self.source 620 } 621 622 pub fn matched_filter(&self) -> &RelayMatchedFilterContext { 623 self.matched_filters 624 .first() 625 .expect("projection context must include at least one matched filter") 626 } 627 628 pub fn matched_filters(&self) -> &[RelayMatchedFilterContext] { 629 &self.matched_filters 630 } 631 632 pub fn event(&self) -> &RelayEventContext { 633 &self.event 634 } 635 636 pub fn authenticated_pubkeys(&self) -> &[PublicKeyHex] { 637 &self.authenticated_pubkeys 638 } 639 } 640 641 impl RelayEventStoredContext { 642 pub fn new(event: RelayEventContext, store_offsets: Vec<u64>) -> Self { 643 Self { 644 event, 645 store_offsets, 646 } 647 } 648 649 pub fn event(&self) -> &RelayEventContext { 650 &self.event 651 } 652 653 pub fn store_offsets(&self) -> &[u64] { 654 &self.store_offsets 655 } 656 } 657 658 struct TanglePocketQueryRateLimitRequest<'a> { 659 scope: TangleRateLimitScope, 660 rules: TangleQueryRateLimitConfig, 661 label: &'static str, 662 subscription_id: &'a SubscriptionId, 663 filters: &'a [PocketOwnedFilter], 664 auth: &'a BaseAuthState, 665 context: TangleClientRateLimitContext, 666 now: UnixTimestamp, 667 } 668 669 #[derive(Debug, Clone, Copy, PartialEq, Eq)] 670 enum TangleQueryClassification { 671 Bounded, 672 Broad(TangleBroadQueryReason), 673 } 674 675 impl TangleQueryClassification { 676 fn is_broad(self) -> bool { 677 matches!(self, Self::Broad(_)) 678 } 679 } 680 681 #[derive(Debug, Clone, Copy, PartialEq, Eq)] 682 enum TangleBroadQueryReason { 683 EmptyFilters, 684 MissingPrimaryConstraint, 685 MissingBoundedSelector, 686 HighLimit, 687 BroadTimeWindow, 688 } 689 690 #[derive(Debug, Clone, Copy)] 691 struct TangleQueryClassifier { 692 limits: BaseRelayLimits, 693 } 694 695 const BROAD_QUERY_TIME_WINDOW_SECONDS: u64 = 31 * 24 * 60 * 60; 696 697 impl TangleQueryClassifier { 698 fn new(limits: BaseRelayLimits) -> Self { 699 Self { limits } 700 } 701 702 fn classify_pocket_query(self, filters: &[PocketOwnedFilter]) -> TangleQueryClassification { 703 if filters.is_empty() { 704 return TangleQueryClassification::Broad(TangleBroadQueryReason::EmptyFilters); 705 } 706 filters 707 .iter() 708 .map(|filter| self.classify_pocket_query_filter(filter)) 709 .find(|classification| classification.is_broad()) 710 .unwrap_or(TangleQueryClassification::Bounded) 711 } 712 713 fn classify_pocket_count(self, filters: &[PocketOwnedFilter]) -> TangleQueryClassification { 714 if filters.is_empty() { 715 return TangleQueryClassification::Broad(TangleBroadQueryReason::EmptyFilters); 716 } 717 filters 718 .iter() 719 .map(|filter| self.classify_pocket_count_filter(filter)) 720 .find(|classification| classification.is_broad()) 721 .unwrap_or(TangleQueryClassification::Bounded) 722 } 723 724 fn classify_pocket_query_filter(self, filter: &PocketFilter) -> TangleQueryClassification { 725 if !self.has_pocket_primary_constraint(filter) { 726 return TangleQueryClassification::Broad( 727 TangleBroadQueryReason::MissingPrimaryConstraint, 728 ); 729 } 730 if self.has_pocket_high_limit(filter) { 731 return TangleQueryClassification::Broad(TangleBroadQueryReason::HighLimit); 732 } 733 if self.has_pocket_broad_time_window(filter) && !self.has_pocket_strong_constraint(filter) { 734 return TangleQueryClassification::Broad(TangleBroadQueryReason::BroadTimeWindow); 735 } 736 TangleQueryClassification::Bounded 737 } 738 739 fn classify_pocket_count_filter(self, filter: &PocketFilter) -> TangleQueryClassification { 740 if !self.has_pocket_primary_constraint(filter) { 741 return TangleQueryClassification::Broad( 742 TangleBroadQueryReason::MissingPrimaryConstraint, 743 ); 744 } 745 if self.has_pocket_high_limit(filter) { 746 return TangleQueryClassification::Broad(TangleBroadQueryReason::HighLimit); 747 } 748 if self.has_pocket_broad_time_window(filter) { 749 return TangleQueryClassification::Broad(TangleBroadQueryReason::BroadTimeWindow); 750 } 751 if !self.has_pocket_count_bounded_selector(filter) { 752 return TangleQueryClassification::Broad( 753 TangleBroadQueryReason::MissingBoundedSelector, 754 ); 755 } 756 TangleQueryClassification::Bounded 757 } 758 759 fn has_pocket_primary_constraint(self, filter: &PocketFilter) -> bool { 760 filter.num_ids() > 0 761 || filter.num_authors() > 0 762 || filter.num_kinds() > 0 763 || self.has_pocket_group_constraint(filter) 764 } 765 766 fn has_pocket_strong_constraint(self, filter: &PocketFilter) -> bool { 767 filter.num_ids() > 0 || filter.num_authors() > 0 || self.has_pocket_group_constraint(filter) 768 } 769 770 fn has_pocket_count_bounded_selector(self, filter: &PocketFilter) -> bool { 771 self.has_pocket_strong_constraint(filter) 772 || (filter.num_kinds() > 0 && self.has_pocket_bounded_time_window(filter)) 773 || self.has_pocket_hll_count_selector(filter) 774 } 775 776 fn has_pocket_hll_count_selector(self, filter: &PocketFilter) -> bool { 777 filter 778 .hyperloglog_offset() 779 .is_ok_and(|offset| offset.is_some()) 780 } 781 782 fn has_pocket_group_constraint(self, filter: &PocketFilter) -> bool { 783 filter 784 .tags() 785 .map(|tags| { 786 tags.iter().any(|mut tag| { 787 let name = tag.next(); 788 let has_value = tag.next().is_some(); 789 matches!(name, Some(value) if matches!(value, b"h" | b"d")) && has_value 790 }) 791 }) 792 .unwrap_or(false) 793 } 794 795 fn has_pocket_high_limit(self, filter: &PocketFilter) -> bool { 796 let limit = if filter.limit() == u32::MAX { 797 self.limits.default_limit() 798 } else { 799 u64::from(filter.limit()) 800 }; 801 limit >= self.limits.max_limit() 802 } 803 804 fn has_pocket_bounded_time_window(self, filter: &PocketFilter) -> bool { 805 if filter.since() == PocketTime::min() || filter.until() == PocketTime::max() { 806 return false; 807 } 808 filter 809 .until() 810 .as_ref() 811 .saturating_sub(*filter.since().as_ref()) 812 <= BROAD_QUERY_TIME_WINDOW_SECONDS 813 } 814 815 fn has_pocket_broad_time_window(self, filter: &PocketFilter) -> bool { 816 if filter.since() == PocketTime::min() || filter.until() == PocketTime::max() { 817 return false; 818 } 819 filter 820 .until() 821 .as_ref() 822 .saturating_sub(*filter.since().as_ref()) 823 > BROAD_QUERY_TIME_WINDOW_SECONDS 824 } 825 } 826 827 impl RelayRuntime { 828 pub fn open(config: BaseRelayRuntimeConfig) -> Result<Self, BaseRelayError> { 829 Self::open_with_hooks(config, Arc::new(NoopRelayRuntimeHooks)) 830 } 831 832 pub fn open_with_hooks( 833 config: BaseRelayRuntimeConfig, 834 hooks: Arc<dyn RelayRuntimeHooks>, 835 ) -> Result<Self, BaseRelayError> { 836 let limits = TangleRuntimeLimits::from_config(&config)?; 837 let relay = config.open_relay()?; 838 let readiness = BaseRelayReadinessHandle::new(relay.readiness_state()); 839 let event_bus = TangleEventBus::new(limits.event_bus_capacity())?; 840 let rate_limiter = TangleRateLimiter::new(); 841 let metrics = TangleRuntimeMetrics::new(); 842 metrics.record_disk_used_bytes(directory_size_bytes( 843 config.pocket_config().data_directory(), 844 )); 845 metrics.record_event_bus_receivers(event_bus.receiver_count()); 846 metrics.record_outbox_pending_events(relay.group_outbox_pending_events()); 847 logging::log_runtime_opened(&config); 848 Ok(Self { 849 config, 850 relay, 851 readiness, 852 event_bus, 853 rate_limiter, 854 metrics, 855 limits, 856 shutdown: TangleShutdownSignal::new(), 857 hooks, 858 }) 859 } 860 861 pub fn config(&self) -> &BaseRelayRuntimeConfig { 862 &self.config 863 } 864 865 pub fn relay(&self) -> &BaseRelay { 866 &self.relay 867 } 868 869 pub fn relay_mut(&mut self) -> &mut BaseRelay { 870 &mut self.relay 871 } 872 873 pub fn auth_state(&self) -> Result<BaseAuthState, BaseRelayError> { 874 self.config.auth_state() 875 } 876 877 pub fn readiness_state(&self) -> BaseRelayReadinessState { 878 self.readiness.snapshot() 879 } 880 881 pub fn readiness_handle(&self) -> BaseRelayReadinessHandle { 882 self.readiness.clone() 883 } 884 885 pub fn limits(&self) -> TangleRuntimeLimits { 886 self.limits 887 } 888 889 pub fn event_bus(&self) -> &TangleEventBus { 890 &self.event_bus 891 } 892 893 pub fn rate_limiter(&self) -> &TangleRateLimiter { 894 &self.rate_limiter 895 } 896 897 pub fn metrics(&self) -> &TangleRuntimeMetrics { 898 &self.metrics 899 } 900 901 pub fn shutdown_signal(&self) -> &TangleShutdownSignal { 902 &self.shutdown 903 } 904 905 pub fn shutdown(&mut self) -> Result<BaseRelayShutdownReport, BaseRelayError> { 906 self.shutdown.request_shutdown(); 907 self.relay.shutdown() 908 } 909 } 910 911 struct RelayRuntimeShared { 912 config: Arc<BaseRelayRuntimeConfig>, 913 store: PocketStoreHandle, 914 groups: Option<GroupServiceHandle>, 915 readiness: BaseRelayReadinessHandle, 916 limits: TangleRuntimeLimits, 917 event_bus: TangleEventBus, 918 rate_limiter: TangleRateLimiter, 919 metrics: TangleRuntimeMetrics, 920 shutdown: TangleShutdownSignal, 921 hooks: Arc<dyn RelayRuntimeHooks>, 922 } 923 924 impl RelayRuntimeShared { 925 fn from_runtime(runtime: RelayRuntime) -> Self { 926 let RelayRuntime { 927 config, 928 relay, 929 readiness, 930 limits, 931 event_bus, 932 rate_limiter, 933 metrics, 934 shutdown, 935 hooks, 936 } = runtime; 937 let store = relay.store_handle(); 938 let groups = relay.group_service_handle(); 939 Self { 940 config: Arc::new(config), 941 store, 942 groups, 943 readiness, 944 limits, 945 event_bus, 946 rate_limiter, 947 metrics, 948 shutdown, 949 hooks, 950 } 951 } 952 953 fn rate_limit_event_pocket( 954 &self, 955 event: &PocketEvent, 956 context: TangleClientRateLimitContext, 957 now: UnixTimestamp, 958 ) -> Result<Option<RelayMessage>, BaseRelayError> { 959 let rules = self.config.rate_limits().event(); 960 if let Some(peer_ip) = context.peer_ip 961 && let Some(message) = self.rate_limit_ok_pocket( 962 event, 963 TangleRateLimitKey::ip(TangleRateLimitScope::Event, peer_ip), 964 rules.per_ip(), 965 "event ip", 966 now, 967 )? 968 { 969 return Ok(Some(message)); 970 } 971 self.rate_limit_ok_pocket( 972 event, 973 TangleRateLimitKey::pubkey(TangleRateLimitScope::Event, pocket_event_pubkey(event)?), 974 rules.per_pubkey(), 975 "event pubkey", 976 now, 977 ) 978 .and_then(|message| { 979 if message.is_some() { 980 return Ok(message); 981 } 982 self.rate_limit_ok_pocket( 983 event, 984 TangleRateLimitKey::kind(TangleRateLimitScope::Event, pocket_event_kind(event)?), 985 rules.per_kind(), 986 "event kind", 987 now, 988 ) 989 }) 990 } 991 992 fn rate_limit_auth_attempt_pocket( 993 &self, 994 event: &PocketEvent, 995 context: TangleClientRateLimitContext, 996 now: UnixTimestamp, 997 ) -> Result<Option<RelayMessage>, BaseRelayError> { 998 let rules = self.config.rate_limits().auth(); 999 if let Some(peer_ip) = context.peer_ip 1000 && let Some(message) = self.rate_limit_ok_pocket( 1001 event, 1002 TangleRateLimitKey::ip(TangleRateLimitScope::Auth, peer_ip), 1003 rules.per_ip(), 1004 "auth ip", 1005 now, 1006 )? 1007 { 1008 return Ok(Some(message)); 1009 } 1010 self.rate_limit_ok_pocket( 1011 event, 1012 TangleRateLimitKey::pubkey(TangleRateLimitScope::Auth, pocket_event_pubkey(event)?), 1013 rules.per_pubkey(), 1014 "auth pubkey", 1015 now, 1016 ) 1017 } 1018 1019 fn rate_limit_auth_failure_pocket( 1020 &self, 1021 event: &PocketEvent, 1022 context: TangleClientRateLimitContext, 1023 now: UnixTimestamp, 1024 ) -> Result<Option<RelayMessage>, BaseRelayError> { 1025 let rules = self.config.rate_limits().auth(); 1026 if let Some(peer_ip) = context.peer_ip 1027 && let Some(message) = self.rate_limit_ok_pocket( 1028 event, 1029 TangleRateLimitKey::auth_failure(Some(peer_ip), None), 1030 rules.failures_per_ip(), 1031 "auth failure ip", 1032 now, 1033 )? 1034 { 1035 return Ok(Some(message)); 1036 } 1037 self.rate_limit_ok_pocket( 1038 event, 1039 TangleRateLimitKey::auth_failure(None, Some(pocket_event_pubkey(event)?)), 1040 rules.failures(), 1041 "auth failure", 1042 now, 1043 ) 1044 } 1045 1046 fn rate_limit_group_write_pocket( 1047 &self, 1048 event: &PocketEvent, 1049 context: TangleClientRateLimitContext, 1050 now: UnixTimestamp, 1051 ) -> Result<Option<RelayMessage>, BaseRelayError> { 1052 if !self.config.groups().enabled() { 1053 return Ok(None); 1054 } 1055 let class = 1056 validate_client_group_event_structure(event, self.config.groups().limits()).ok(); 1057 let Some(class) = class else { 1058 return Ok(None); 1059 }; 1060 let Some(group_id) = class.group_id().cloned() else { 1061 return Ok(None); 1062 }; 1063 let rules = self.config.rate_limits().group(); 1064 let kind = pocket_event_kind(event)?; 1065 let pubkey = pocket_event_pubkey(event)?; 1066 if kind.as_u32() == KIND_GROUP_JOIN_REQUEST { 1067 if let Some(peer_ip) = context.peer_ip 1068 && let Some(message) = self.rate_limit_ok_pocket( 1069 event, 1070 TangleRateLimitKey::join_flow_ip(group_id.clone(), peer_ip), 1071 rules.join_flow_per_ip(), 1072 "group join ip", 1073 now, 1074 )? 1075 { 1076 return Ok(Some(message)); 1077 } 1078 if let Some(message) = self.rate_limit_ok_pocket( 1079 event, 1080 TangleRateLimitKey::join_flow(group_id.clone(), pubkey.clone()), 1081 rules.join_flow(), 1082 "group join", 1083 now, 1084 )? { 1085 return Ok(Some(message)); 1086 } 1087 } 1088 if let Some(peer_ip) = context.peer_ip 1089 && let Some(message) = self.rate_limit_ok_pocket( 1090 event, 1091 TangleRateLimitKey::ip(TangleRateLimitScope::GroupWrite, peer_ip), 1092 rules.write_per_ip(), 1093 "group ip", 1094 now, 1095 )? 1096 { 1097 return Ok(Some(message)); 1098 } 1099 if let Some(message) = self.rate_limit_ok_pocket( 1100 event, 1101 TangleRateLimitKey::pubkey(TangleRateLimitScope::GroupWrite, pubkey), 1102 rules.write_per_pubkey(), 1103 "group pubkey", 1104 now, 1105 )? { 1106 return Ok(Some(message)); 1107 } 1108 if let Some(message) = self.rate_limit_ok_pocket( 1109 event, 1110 TangleRateLimitKey::group(TangleRateLimitScope::GroupWrite, group_id), 1111 rules.write_per_group(), 1112 "group write", 1113 now, 1114 )? { 1115 return Ok(Some(message)); 1116 } 1117 self.rate_limit_ok_pocket( 1118 event, 1119 TangleRateLimitKey::kind(TangleRateLimitScope::GroupWrite, kind), 1120 rules.write_per_kind(), 1121 "group kind", 1122 now, 1123 ) 1124 } 1125 1126 fn is_group_event_pocket(&self, event: &PocketEvent) -> bool { 1127 self.config.groups().enabled() 1128 && validate_client_group_event_structure(event, self.config.groups().limits()) 1129 .is_ok_and(|class| !matches!(class, GroupEventClass::NonGroup)) 1130 } 1131 1132 fn handle_pocket_event_with_auth_report( 1133 &self, 1134 event: &PocketEvent, 1135 auth: &BaseAuthState, 1136 ) -> Result<BaseRelayEventWrite, BaseRelayError> { 1137 BaseRelay::handle_pocket_event_with_shared_services( 1138 &self.store, 1139 self.groups.as_ref(), 1140 self.limits.base_relay_limits(), 1141 event, 1142 auth, 1143 ) 1144 } 1145 1146 fn group_outbox_pending_events(&self) -> usize { 1147 self.groups 1148 .as_ref() 1149 .map(GroupServiceHandle::outbox_pending_events) 1150 .unwrap_or(0) 1151 } 1152 1153 fn query_req_with_auth_report( 1154 &self, 1155 subscription_id: SubscriptionId, 1156 filters: Vec<PocketOwnedFilter>, 1157 search_present: bool, 1158 auth: &BaseAuthState, 1159 ) -> Result<BaseRelayQueryReport, BaseRelayError> { 1160 BaseRelay::query_req_with_shared_services( 1161 &self.store, 1162 self.groups.as_ref(), 1163 self.limits.base_relay_limits(), 1164 self.config.pocket_query_config(), 1165 BaseRelayReqQuery::new(subscription_id, filters, search_present, auth), 1166 ) 1167 } 1168 1169 fn project_query_report( 1170 &self, 1171 report: BaseRelayQueryReport, 1172 filters: &[PocketOwnedFilter], 1173 projection: &RelayProjectionContext, 1174 auth: &BaseAuthState, 1175 ) -> Result<BaseRelayQueryReport, BaseRelayError> { 1176 let group_read_denied = report.group_read_denied(); 1177 let query_metrics = report.query_metrics(); 1178 let messages = 1179 self.project_runtime_messages(report.into_messages(), filters, projection, auth)?; 1180 let returned_events = messages 1181 .iter() 1182 .filter(|message| matches!(message, RuntimeRelayMessage::Event { .. })) 1183 .count(); 1184 let query_metrics = query_metrics.with_returned_events(returned_events); 1185 Ok(BaseRelayQueryReport::new( 1186 messages, 1187 group_read_denied, 1188 query_metrics, 1189 )) 1190 } 1191 1192 fn project_runtime_messages( 1193 &self, 1194 messages: Vec<RuntimeRelayMessage>, 1195 filters: &[PocketOwnedFilter], 1196 projection: &RelayProjectionContext, 1197 auth: &BaseAuthState, 1198 ) -> Result<Vec<RuntimeRelayMessage>, BaseRelayError> { 1199 let mut output = Vec::with_capacity(messages.len()); 1200 let mut event_ids = BTreeSet::new(); 1201 for message in messages { 1202 match message { 1203 RuntimeRelayMessage::Event { 1204 subscription_id, 1205 event, 1206 } => { 1207 let matched_filters = self.matched_filters_for_event(filters, &event)?; 1208 if let Some(projected) = 1209 self.project_event_output(RelayEventProjectionRequest { 1210 subscription_id: &subscription_id, 1211 projection, 1212 source: RelayEventProjectionSource::HistoricalQuery, 1213 event: &event, 1214 auth, 1215 matched_filters, 1216 })? 1217 && event_ids.insert(projected.id()) 1218 { 1219 output.push(RuntimeRelayMessage::event(subscription_id, projected)); 1220 } 1221 } 1222 message => output.push(message), 1223 } 1224 } 1225 Ok(output) 1226 } 1227 1228 fn project_event_output( 1229 &self, 1230 request: RelayEventProjectionRequest<'_>, 1231 ) -> Result<Option<PocketOwnedEvent>, BaseRelayError> { 1232 let RelayEventProjectionRequest { 1233 subscription_id, 1234 projection, 1235 source, 1236 event, 1237 auth, 1238 matched_filters, 1239 } = request; 1240 let context = 1241 RelayEventProjectionContext::new_with_matched_filters_and_authenticated_pubkeys( 1242 subscription_id.clone(), 1243 projection.clone(), 1244 source, 1245 matched_filters 1246 .iter() 1247 .map(|(matched_filter, _)| RelayMatchedFilterContext::from_base(matched_filter)) 1248 .collect(), 1249 RelayEventContext::from_pocket_event(event)?, 1250 auth.authenticated_pubkeys().iter().cloned().collect(), 1251 ); 1252 match self.hooks.project_event(&context) { 1253 RelayEventProjectionDecision::Emit => Ok(Some(event.to_owned())), 1254 RelayEventProjectionDecision::Suppress => Ok(None), 1255 RelayEventProjectionDecision::ReplaceWithStoredOffset { store_offset } => { 1256 let Ok(replacement) = self.store.event_by_offset(store_offset) else { 1257 return Ok(None); 1258 }; 1259 let group_auth = 1260 GroupAuthContext::new(auth.authenticated_pubkeys().iter().cloned()); 1261 if BaseRelay::group_read_gate_visible_to_auth( 1262 self.groups.as_ref(), 1263 &replacement, 1264 &group_auth, 1265 )? && matched_filters 1266 .iter() 1267 .try_fold(false, |matched, (_, filter)| { 1268 if matched { 1269 return Ok(true); 1270 } 1271 filter 1272 .event_matches(&replacement) 1273 .map_err(|error| BaseRelayError::error(error.to_string())) 1274 })? 1275 { 1276 Ok(Some(replacement)) 1277 } else { 1278 Ok(None) 1279 } 1280 } 1281 } 1282 } 1283 1284 fn matched_filters_for_event<'a>( 1285 &self, 1286 filters: &'a [PocketOwnedFilter], 1287 event: &PocketEvent, 1288 ) -> Result<Vec<(BaseRelayMatchedFilterContext, &'a PocketFilter)>, BaseRelayError> { 1289 let mut matched_filters = Vec::new(); 1290 for (filter_index, filter) in filters.iter().enumerate() { 1291 if filter 1292 .event_matches(event) 1293 .map_err(|error| BaseRelayError::error(error.to_string()))? 1294 { 1295 let filter: &PocketFilter = filter; 1296 matched_filters.push((matched_filter_context(filter_index, filter), filter)); 1297 } 1298 } 1299 if matched_filters.is_empty() { 1300 return Err(BaseRelayError::error( 1301 "query output did not match any request filter", 1302 )); 1303 } 1304 Ok(matched_filters) 1305 } 1306 1307 fn query_projected_req_with_auth_report( 1308 &self, 1309 subscription_id: SubscriptionId, 1310 filters: Vec<PocketOwnedFilter>, 1311 search_present: bool, 1312 auth: &BaseAuthState, 1313 projection: &RelayProjectionContext, 1314 candidate_limit: NonZeroU32, 1315 ) -> Result<BaseRelayQueryReport, BaseRelayError> { 1316 self.limits 1317 .base_relay_limits() 1318 .validate_subscription_id(&subscription_id)?; 1319 self.limits 1320 .base_relay_limits() 1321 .validate_pocket_filters(&filters)?; 1322 if let Some(message) = 1323 BaseRelay::unsupported_search_present_closed(&subscription_id, search_present) 1324 { 1325 return Ok(BaseRelayQueryReport::new( 1326 vec![message.into()], 1327 false, 1328 BaseRelayQueryMetrics::default(), 1329 )); 1330 } 1331 let group_auth = GroupAuthContext::new(auth.authenticated_pubkeys().iter().cloned()); 1332 let mut output = Vec::new(); 1333 let mut group_read_denied = false; 1334 let mut query_metrics = BaseRelayQueryMetrics::default(); 1335 for filter in &filters { 1336 let report = BaseRelay::query_filter_events_report_with_services( 1337 &self.store, 1338 self.groups.as_ref(), 1339 self.limits.base_relay_limits(), 1340 self.config.pocket_query_config(), 1341 filter, 1342 &group_auth, 1343 BaseRelayFilterLimitMode::Override(candidate_limit.get()), 1344 )?; 1345 group_read_denied |= report.group_read_denied(); 1346 query_metrics = query_metrics.add(report.query_metrics()); 1347 let events = BaseRelay::sort_and_dedupe_query_events(report.into_events()); 1348 let mut projected = Vec::new(); 1349 for event in events { 1350 let matched_filters = self.matched_filters_for_event(&filters, &event)?; 1351 if let Some(event) = self.project_event_output(RelayEventProjectionRequest { 1352 subscription_id: &subscription_id, 1353 projection, 1354 source: RelayEventProjectionSource::HistoricalQuery, 1355 event: &event, 1356 auth, 1357 matched_filters, 1358 })? { 1359 projected.push(event); 1360 } 1361 } 1362 let mut projected = BaseRelay::sort_and_dedupe_query_events(projected); 1363 projected.truncate( 1364 self.limits 1365 .base_relay_limits() 1366 .effective_pocket_filter_limit_for_query(filter), 1367 ); 1368 output.extend(projected); 1369 } 1370 let events = BaseRelay::sort_and_dedupe_query_events(output); 1371 let query_metrics = query_metrics.with_returned_events(events.len()); 1372 let mut messages = events 1373 .into_iter() 1374 .map(|event| RuntimeRelayMessage::event(subscription_id.clone(), event)) 1375 .collect::<Vec<_>>(); 1376 if group_read_denied { 1377 let group_auth = GroupAuthContext::new(auth.authenticated_pubkeys().iter().cloned()); 1378 messages.push(BaseRelay::redacted_req_closed(subscription_id, &group_auth).into()); 1379 } else { 1380 messages.push(RelayMessage::Eose(subscription_id).into()); 1381 } 1382 Ok(BaseRelayQueryReport::new( 1383 messages, 1384 group_read_denied, 1385 query_metrics, 1386 )) 1387 } 1388 1389 fn handle_count_with_auth_report( 1390 &self, 1391 subscription_id: SubscriptionId, 1392 filters: Vec<PocketOwnedFilter>, 1393 search_present: bool, 1394 auth: &BaseAuthState, 1395 ) -> Result<BaseRelayCountReport, BaseRelayError> { 1396 BaseRelay::handle_count_with_shared_services( 1397 &self.store, 1398 self.groups.as_ref(), 1399 self.limits.base_relay_limits(), 1400 self.config.pocket_query_config(), 1401 BaseRelayCountQuery::new(subscription_id, filters, search_present, auth), 1402 ) 1403 } 1404 1405 fn rate_limit_req_pocket( 1406 &self, 1407 subscription_id: &SubscriptionId, 1408 filters: &[PocketOwnedFilter], 1409 auth: &BaseAuthState, 1410 context: TangleClientRateLimitContext, 1411 now: UnixTimestamp, 1412 ) -> Option<RelayMessage> { 1413 self.rate_limit_pocket_query(TanglePocketQueryRateLimitRequest { 1414 scope: TangleRateLimitScope::Req, 1415 rules: self.config.rate_limits().req(), 1416 label: "req", 1417 subscription_id, 1418 filters, 1419 auth, 1420 context, 1421 now, 1422 }) 1423 } 1424 1425 fn rate_limit_count_pocket( 1426 &self, 1427 subscription_id: &SubscriptionId, 1428 filters: &[PocketOwnedFilter], 1429 auth: &BaseAuthState, 1430 context: TangleClientRateLimitContext, 1431 now: UnixTimestamp, 1432 ) -> Option<RelayMessage> { 1433 self.rate_limit_pocket_query(TanglePocketQueryRateLimitRequest { 1434 scope: TangleRateLimitScope::Count, 1435 rules: self.config.rate_limits().count(), 1436 label: "count", 1437 subscription_id, 1438 filters, 1439 auth, 1440 context, 1441 now, 1442 }) 1443 } 1444 1445 fn refuse_broad_count( 1446 &self, 1447 subscription_id: &SubscriptionId, 1448 filters: &[PocketOwnedFilter], 1449 ) -> Option<RelayMessage> { 1450 if TangleQueryClassifier::new(self.limits.base_relay_limits()) 1451 .classify_pocket_count(filters) 1452 .is_broad() 1453 { 1454 self.metrics.record_count_refusal(); 1455 self.metrics.record_broad_query_rejection(); 1456 return Some(RelayMessage::Closed { 1457 subscription_id: subscription_id.clone(), 1458 message: BaseRelayError::restricted("count filters are too broad or expensive") 1459 .prefixed_message(), 1460 }); 1461 } 1462 None 1463 } 1464 1465 fn rate_limit_pocket_query( 1466 &self, 1467 request: TanglePocketQueryRateLimitRequest<'_>, 1468 ) -> Option<RelayMessage> { 1469 if let Some(peer_ip) = request.context.peer_ip 1470 && let Some(message) = self.rate_limit_closed( 1471 request.subscription_id, 1472 TangleRateLimitKey::ip(request.scope, peer_ip), 1473 request.rules.per_ip(), 1474 request.label, 1475 "ip", 1476 request.now, 1477 ) 1478 { 1479 return Some(message); 1480 } 1481 if let Some(connection_id) = request.context.connection_id 1482 && let Some(message) = self.rate_limit_closed( 1483 request.subscription_id, 1484 TangleRateLimitKey::connection(request.scope, connection_id), 1485 request.rules.per_connection(), 1486 request.label, 1487 "connection", 1488 request.now, 1489 ) 1490 { 1491 return Some(message); 1492 } 1493 for pubkey in request.auth.authenticated_pubkeys() { 1494 if let Some(message) = self.rate_limit_closed( 1495 request.subscription_id, 1496 TangleRateLimitKey::pubkey(request.scope, pubkey.clone()), 1497 request.rules.per_pubkey(), 1498 request.label, 1499 "pubkey", 1500 request.now, 1501 ) { 1502 return Some(message); 1503 } 1504 } 1505 for group_id in pocket_filter_group_ids(request.filters) { 1506 if let Some(message) = self.rate_limit_closed( 1507 request.subscription_id, 1508 TangleRateLimitKey::group(request.scope, group_id), 1509 request.rules.per_group(), 1510 request.label, 1511 "group", 1512 request.now, 1513 ) { 1514 return Some(message); 1515 } 1516 } 1517 for kind in pocket_filter_kinds(request.filters) { 1518 if let Some(message) = self.rate_limit_closed( 1519 request.subscription_id, 1520 TangleRateLimitKey::kind(request.scope, kind), 1521 request.rules.per_kind(), 1522 request.label, 1523 "kind", 1524 request.now, 1525 ) { 1526 return Some(message); 1527 } 1528 } 1529 let classifier = TangleQueryClassifier::new(self.limits.base_relay_limits()); 1530 let query_classification = match request.scope { 1531 TangleRateLimitScope::Req => classifier.classify_pocket_query(request.filters), 1532 TangleRateLimitScope::Count => classifier.classify_pocket_count(request.filters), 1533 TangleRateLimitScope::Auth 1534 | TangleRateLimitScope::Event 1535 | TangleRateLimitScope::GroupWrite => classifier.classify_pocket_query(request.filters), 1536 }; 1537 if query_classification.is_broad() 1538 && let Some(message) = self.rate_limit_closed( 1539 request.subscription_id, 1540 TangleRateLimitKey::query_class(request.scope, TangleRateLimitQueryClass::Broad), 1541 request.rules.broad(), 1542 request.label, 1543 "broad", 1544 request.now, 1545 ) 1546 { 1547 self.metrics.record_broad_query_rejection(); 1548 return Some(message); 1549 } 1550 None 1551 } 1552 1553 fn rate_limit_closed( 1554 &self, 1555 subscription_id: &SubscriptionId, 1556 key: TangleRateLimitKey, 1557 rule: TangleRateLimitRule, 1558 label: &'static str, 1559 dimension: &'static str, 1560 now: UnixTimestamp, 1561 ) -> Option<RelayMessage> { 1562 match self.rate_limiter.record(key, rule, now) { 1563 TangleRateLimitDecision::Allowed { .. } => None, 1564 TangleRateLimitDecision::Rejected { reset_at } => { 1565 self.metrics.record_rate_limit_rejection(); 1566 logging::log_rate_limit_rejected(label, dimension, reset_at); 1567 Some(RelayMessage::Closed { 1568 subscription_id: subscription_id.clone(), 1569 message: BaseRelayError::rate_limited(format!( 1570 "{label} {dimension} rate limit exceeded until {reset_at}" 1571 )) 1572 .prefixed_message(), 1573 }) 1574 } 1575 } 1576 } 1577 1578 fn rate_limit_ok_pocket( 1579 &self, 1580 event: &PocketEvent, 1581 key: TangleRateLimitKey, 1582 rule: TangleRateLimitRule, 1583 label: &'static str, 1584 now: UnixTimestamp, 1585 ) -> Result<Option<RelayMessage>, BaseRelayError> { 1586 Ok(match self.rate_limiter.record(key, rule, now) { 1587 TangleRateLimitDecision::Allowed { .. } => None, 1588 TangleRateLimitDecision::Rejected { reset_at } => { 1589 self.metrics.record_rate_limit_rejection(); 1590 logging::log_rate_limit_rejected(label, "event", reset_at); 1591 Some(RelayMessage::Ok { 1592 event_id: pocket_event_id(event)?, 1593 accepted: false, 1594 message: BaseRelayError::rate_limited(format!( 1595 "{label} rate limit exceeded until {reset_at}" 1596 )) 1597 .prefixed_message(), 1598 }) 1599 } 1600 }) 1601 } 1602 } 1603 1604 #[derive(Clone)] 1605 pub struct RelayRuntimeHandle { 1606 inner: Arc<RelayRuntimeShared>, 1607 } 1608 1609 impl RelayRuntimeHandle { 1610 pub fn new(runtime: RelayRuntime) -> Self { 1611 Self { 1612 inner: Arc::new(RelayRuntimeShared::from_runtime(runtime)), 1613 } 1614 } 1615 1616 pub fn metrics(&self) -> TangleRuntimeMetrics { 1617 self.inner.metrics.clone() 1618 } 1619 1620 pub fn readiness_handle(&self) -> BaseRelayReadinessHandle { 1621 self.inner.readiness.clone() 1622 } 1623 1624 pub(crate) fn sanitize_public_message( 1625 &self, 1626 message: RuntimeRelayMessage, 1627 ) -> RuntimeRelayMessage { 1628 message.map_protocol(|message| self.inner.hooks.sanitize_public_message(message)) 1629 } 1630 1631 pub fn limits(&self) -> TangleRuntimeLimits { 1632 self.inner.limits 1633 } 1634 1635 pub async fn auth_state(&self) -> Result<BaseAuthState, BaseRelayError> { 1636 self.inner.config.auth_state() 1637 } 1638 1639 pub async fn handle_count_pocket( 1640 &self, 1641 subscription_id: SubscriptionId, 1642 filters: Vec<PocketOwnedFilter>, 1643 auth: &mut BaseAuthState, 1644 now: UnixTimestamp, 1645 ) -> Result<Vec<RelayMessage>, BaseRelayError> { 1646 let messages = self 1647 .handle_client_message_with_rate_limit_context( 1648 RuntimeClientMessage::Count { 1649 subscription_id, 1650 filters, 1651 search_present: false, 1652 }, 1653 auth, 1654 TangleClientRateLimitContext::default(), 1655 now, 1656 ) 1657 .await?; 1658 protocol_control_messages(messages) 1659 } 1660 1661 #[cfg(test)] 1662 pub(crate) async fn handle_client_message( 1663 &self, 1664 message: RuntimeClientMessage, 1665 auth: &mut BaseAuthState, 1666 now: UnixTimestamp, 1667 ) -> Result<Vec<RelayMessage>, BaseRelayError> { 1668 let messages = self 1669 .handle_client_message_with_rate_limit_context( 1670 message, 1671 auth, 1672 TangleClientRateLimitContext::default(), 1673 now, 1674 ) 1675 .await?; 1676 protocol_messages_for_test(messages) 1677 } 1678 1679 #[cfg(test)] 1680 pub(crate) async fn handle_protocol_client_message_for_test( 1681 &self, 1682 message: tangle_protocol::ClientMessage, 1683 auth: &mut BaseAuthState, 1684 now: UnixTimestamp, 1685 ) -> Result<Vec<RelayMessage>, BaseRelayError> { 1686 self.handle_client_message( 1687 protocol_client_message_to_runtime_for_test(message)?, 1688 auth, 1689 now, 1690 ) 1691 .await 1692 } 1693 1694 #[cfg(test)] 1695 pub(crate) async fn handle_protocol_client_message_with_rate_limit_context_for_test( 1696 &self, 1697 message: tangle_protocol::ClientMessage, 1698 auth: &mut BaseAuthState, 1699 rate_limit_context: TangleClientRateLimitContext, 1700 now: UnixTimestamp, 1701 ) -> Result<Vec<RelayMessage>, BaseRelayError> { 1702 let messages = self 1703 .handle_client_message_with_rate_limit_context( 1704 protocol_client_message_to_runtime_for_test(message)?, 1705 auth, 1706 rate_limit_context, 1707 now, 1708 ) 1709 .await?; 1710 protocol_messages_for_test(messages) 1711 } 1712 1713 pub(crate) async fn handle_client_message_with_rate_limit_context( 1714 &self, 1715 message: RuntimeClientMessage, 1716 auth: &mut BaseAuthState, 1717 rate_limit_context: TangleClientRateLimitContext, 1718 now: UnixTimestamp, 1719 ) -> Result<Vec<RuntimeRelayMessage>, BaseRelayError> { 1720 self.inner 1721 .metrics 1722 .record_client_message(runtime_client_message_metric_kind(&message)); 1723 match message { 1724 RuntimeClientMessage::Event(pocket_event) => { 1725 let started_at = Instant::now(); 1726 let event_id = pocket_event_id(&pocket_event)?; 1727 let event_context = RelayEventContext::from_pocket_event(&pocket_event)?; 1728 let is_group_event = self.inner.is_group_event_pocket(&pocket_event); 1729 if let Some(message) = 1730 self.inner 1731 .rate_limit_event_pocket(&pocket_event, rate_limit_context, now)? 1732 { 1733 record_event_metrics(&self.inner.metrics, &message, is_group_event, started_at); 1734 return Ok(vec![message.into()]); 1735 } 1736 if let Some(message) = self.inner.rate_limit_group_write_pocket( 1737 &pocket_event, 1738 rate_limit_context, 1739 now, 1740 )? { 1741 record_event_metrics(&self.inner.metrics, &message, is_group_event, started_at); 1742 return Ok(vec![message.into()]); 1743 } 1744 let authenticated_pubkeys = auth 1745 .authenticated_pubkeys() 1746 .iter() 1747 .map(ToString::to_string) 1748 .collect(); 1749 let admission = RelayEventAdmissionContext::new( 1750 event_context.clone(), 1751 authenticated_pubkeys, 1752 rate_limit_context.peer_ip(), 1753 rate_limit_context.connection_id(), 1754 now.as_u64(), 1755 ); 1756 let admission_decision = self.inner.hooks.admit_event(&admission); 1757 if let Some(used_bytes) = self.inner.hooks.storage_used_bytes() { 1758 self.inner.metrics.record_disk_used_bytes(used_bytes); 1759 } 1760 if let EventAdmissionDecision::Reject { message } = admission_decision { 1761 let message = RelayMessage::Ok { 1762 event_id, 1763 accepted: false, 1764 message: BaseRelayError::restricted(message).prefixed_message(), 1765 }; 1766 record_event_metrics(&self.inner.metrics, &message, is_group_event, started_at); 1767 return Ok(vec![message.into()]); 1768 } 1769 let result = self 1770 .inner 1771 .handle_pocket_event_with_auth_report(&pocket_event, auth)?; 1772 let group_outbox_pending_events = 1773 is_group_event.then(|| self.inner.group_outbox_pending_events()); 1774 if is_group_event { 1775 for _ in 0..result.stored_offsets().len().saturating_sub(1) { 1776 self.inner.metrics.record_outbox_replayed_event(); 1777 } 1778 self.inner 1779 .metrics 1780 .record_outbox_pending_events(group_outbox_pending_events.unwrap_or(0)); 1781 } 1782 if !result.stored_offsets().is_empty() { 1783 logging::log_event_stored( 1784 &event_id, 1785 result.stored_offsets().len(), 1786 self.inner.metrics.stored_event_offsets(), 1787 ); 1788 self.inner.hooks.event_stored(&RelayEventStoredContext::new( 1789 event_context, 1790 result 1791 .stored_offsets() 1792 .iter() 1793 .map(|offset| offset.as_u64()) 1794 .collect(), 1795 )); 1796 if let Some(used_bytes) = self.inner.hooks.storage_used_bytes() { 1797 self.inner.metrics.record_disk_used_bytes(used_bytes); 1798 } 1799 } 1800 for offset in result.stored_offsets() { 1801 self.inner.metrics.record_stored_event_offset(); 1802 let receivers = self.inner.event_bus.publish(*offset); 1803 self.inner.metrics.record_event_bus_publish(receivers); 1804 } 1805 let message = result.into_message(); 1806 record_event_metrics(&self.inner.metrics, &message, is_group_event, started_at); 1807 Ok(vec![message.into()]) 1808 } 1809 RuntimeClientMessage::Req { 1810 subscription_id, 1811 filters, 1812 search_present, 1813 } => { 1814 let started_at = Instant::now(); 1815 self.inner 1816 .limits 1817 .base_relay_limits() 1818 .validate_subscription_id(&subscription_id)?; 1819 self.inner 1820 .limits 1821 .base_relay_limits() 1822 .validate_pocket_filters(&filters)?; 1823 if let Some(message) = 1824 BaseRelay::unsupported_search_present_closed(&subscription_id, search_present) 1825 { 1826 self.inner 1827 .metrics 1828 .record_query_latency(elapsed_micros(started_at)); 1829 return Ok(vec![message.into()]); 1830 } 1831 if let Some(message) = self.inner.rate_limit_req_pocket( 1832 &subscription_id, 1833 &filters, 1834 auth, 1835 rate_limit_context, 1836 now, 1837 ) { 1838 self.inner 1839 .metrics 1840 .record_query_latency(elapsed_micros(started_at)); 1841 return Ok(vec![message.into()]); 1842 } 1843 let report = self.inner.query_req_with_auth_report( 1844 subscription_id, 1845 filters.clone(), 1846 search_present, 1847 auth, 1848 )?; 1849 let report = self.inner.project_query_report( 1850 report, 1851 &filters, 1852 &RelayProjectionContext::default(), 1853 auth, 1854 )?; 1855 self.inner 1856 .metrics 1857 .record_query_metrics(report.query_metrics()); 1858 if report.group_read_denied() { 1859 self.inner.metrics.record_group_read_denial(); 1860 } 1861 self.inner 1862 .metrics 1863 .record_query_latency(elapsed_micros(started_at)); 1864 Ok(report.into_messages()) 1865 } 1866 RuntimeClientMessage::Count { 1867 subscription_id, 1868 filters, 1869 search_present, 1870 } => { 1871 let started_at = Instant::now(); 1872 self.inner 1873 .limits 1874 .base_relay_limits() 1875 .validate_subscription_id(&subscription_id)?; 1876 self.inner 1877 .limits 1878 .base_relay_limits() 1879 .validate_pocket_filters(&filters)?; 1880 if let Some(message) = 1881 BaseRelay::unsupported_search_present_closed(&subscription_id, search_present) 1882 { 1883 self.inner 1884 .metrics 1885 .record_query_latency(elapsed_micros(started_at)); 1886 return Ok(vec![message.into()]); 1887 } 1888 if let Some(message) = self.inner.refuse_broad_count(&subscription_id, &filters) { 1889 self.inner 1890 .metrics 1891 .record_query_latency(elapsed_micros(started_at)); 1892 return Ok(vec![message.into()]); 1893 } 1894 if let Some(message) = self.inner.rate_limit_count_pocket( 1895 &subscription_id, 1896 &filters, 1897 auth, 1898 rate_limit_context, 1899 now, 1900 ) { 1901 self.inner 1902 .metrics 1903 .record_query_latency(elapsed_micros(started_at)); 1904 return Ok(vec![message.into()]); 1905 } 1906 let report = self.inner.handle_count_with_auth_report( 1907 subscription_id, 1908 filters, 1909 search_present, 1910 auth, 1911 )?; 1912 self.inner 1913 .metrics 1914 .record_query_metrics(report.query_metrics()); 1915 if report.group_read_denied() { 1916 self.inner.metrics.record_group_read_denial(); 1917 } 1918 self.inner 1919 .metrics 1920 .record_query_latency(elapsed_micros(started_at)); 1921 Ok(vec![report.into_message().into()]) 1922 } 1923 RuntimeClientMessage::Auth(pocket_event) => { 1924 let event_id = pocket_event_id(&pocket_event)?; 1925 if let Err(error) = self 1926 .inner 1927 .limits 1928 .base_relay_limits() 1929 .validate_pocket_event(&pocket_event) 1930 { 1931 self.inner.metrics.record_auth_failure(); 1932 return Ok(vec![RuntimeRelayMessage::from(RelayMessage::Ok { 1933 event_id, 1934 accepted: false, 1935 message: error.prefixed_message(), 1936 })]); 1937 } 1938 if let Some(message) = self.inner.rate_limit_auth_attempt_pocket( 1939 &pocket_event, 1940 rate_limit_context, 1941 now, 1942 )? { 1943 self.inner.metrics.record_auth_failure(); 1944 return Ok(vec![message.into()]); 1945 } 1946 let event_for_failure = pocket_event.clone(); 1947 let replies = BaseRelay::handle_pocket_auth_with_limits( 1948 self.inner.limits.base_relay_limits(), 1949 &pocket_event, 1950 auth, 1951 now, 1952 ); 1953 if auth_response_failed(&replies) { 1954 self.inner.metrics.record_auth_failure(); 1955 if let Some(message) = self.inner.rate_limit_auth_failure_pocket( 1956 &event_for_failure, 1957 rate_limit_context, 1958 now, 1959 )? { 1960 return Ok(vec![message.into()]); 1961 } 1962 } else { 1963 self.inner.metrics.record_auth_success(); 1964 } 1965 Ok(replies.into_iter().map(Into::into).collect()) 1966 } 1967 RuntimeClientMessage::Close(subscription_id) => { 1968 self.inner 1969 .limits 1970 .base_relay_limits() 1971 .validate_subscription_id(&subscription_id)?; 1972 Ok(Vec::new()) 1973 } 1974 RuntimeClientMessage::NegOpen { 1975 subscription_id, .. 1976 } 1977 | RuntimeClientMessage::NegMsg { 1978 subscription_id, .. 1979 } => { 1980 self.inner 1981 .limits 1982 .base_relay_limits() 1983 .validate_subscription_id(&subscription_id)?; 1984 Ok(vec![ 1985 BaseRelay::disabled_negentropy_message(subscription_id).into(), 1986 ]) 1987 } 1988 RuntimeClientMessage::NegClose(subscription_id) => { 1989 self.inner 1990 .limits 1991 .base_relay_limits() 1992 .validate_subscription_id(&subscription_id)?; 1993 Ok(Vec::new()) 1994 } 1995 } 1996 } 1997 1998 pub async fn subscribe_events(&self) -> TangleEventReceiver { 1999 let receiver = self.inner.event_bus.subscribe(); 2000 self.inner 2001 .metrics 2002 .record_event_bus_receivers(self.inner.event_bus.receiver_count()); 2003 receiver 2004 } 2005 2006 pub async fn rate_limiter(&self) -> TangleRateLimiter { 2007 self.inner.rate_limiter.clone() 2008 } 2009 2010 pub(crate) async fn rate_limit_req_pocket( 2011 &self, 2012 subscription_id: &SubscriptionId, 2013 filters: &[PocketOwnedFilter], 2014 auth: &BaseAuthState, 2015 rate_limit_context: TangleClientRateLimitContext, 2016 now: UnixTimestamp, 2017 ) -> Option<RelayMessage> { 2018 self.inner 2019 .rate_limit_req_pocket(subscription_id, filters, auth, rate_limit_context, now) 2020 } 2021 2022 pub(crate) async fn query_req_with_auth_report_with_projection_context( 2023 &self, 2024 subscription_id: SubscriptionId, 2025 filters: Vec<PocketOwnedFilter>, 2026 search_present: bool, 2027 auth: &BaseAuthState, 2028 projection: &RelayProjectionContext, 2029 ) -> Result<BaseRelayQueryReport, BaseRelayError> { 2030 let started_at = Instant::now(); 2031 let context = RelayQueryProjectionContext::new_with_authenticated_pubkeys( 2032 subscription_id.clone(), 2033 projection.clone(), 2034 filters 2035 .iter() 2036 .enumerate() 2037 .map(|(index, filter)| { 2038 RelayMatchedFilterContext::from_base(&matched_filter_context(index, filter)) 2039 }) 2040 .collect(), 2041 auth.authenticated_pubkeys().iter().cloned().collect(), 2042 ); 2043 let plan = self.inner.hooks.plan_query(&context); 2044 let report = match plan.limit() { 2045 RelayProjectionQueryLimit::BeforeProjection => { 2046 let report = self.inner.query_req_with_auth_report( 2047 subscription_id, 2048 filters.clone(), 2049 search_present, 2050 auth, 2051 )?; 2052 self.inner 2053 .project_query_report(report, &filters, projection, auth)? 2054 } 2055 RelayProjectionQueryLimit::AfterProjection { candidate_limit } => { 2056 self.inner.query_projected_req_with_auth_report( 2057 subscription_id, 2058 filters, 2059 search_present, 2060 auth, 2061 projection, 2062 candidate_limit, 2063 )? 2064 } 2065 }; 2066 if report.group_read_denied() { 2067 self.inner.metrics.record_group_read_denial(); 2068 } 2069 self.inner 2070 .metrics 2071 .record_query_latency(elapsed_micros(started_at)); 2072 Ok(report) 2073 } 2074 2075 pub async fn event_by_offset_with_auth( 2076 &self, 2077 offset: StoreOffset, 2078 auth: &BaseAuthState, 2079 ) -> Result<Option<PocketOwnedEvent>, BaseRelayError> { 2080 let pocket_event = self.inner.store.event_by_offset(offset.as_u64())?; 2081 let group_auth = GroupAuthContext::new(auth.authenticated_pubkeys().iter().cloned()); 2082 let visible = BaseRelay::group_read_gate_visible_to_auth( 2083 self.inner.groups.as_ref(), 2084 &pocket_event, 2085 &group_auth, 2086 )?; 2087 if !visible { 2088 self.inner.metrics.record_group_read_denial(); 2089 return Ok(None); 2090 } 2091 Ok(Some(pocket_event)) 2092 } 2093 2094 pub(crate) async fn fanout_event_offset_with_projection_context( 2095 &self, 2096 offset: StoreOffset, 2097 subscriptions: &mut LiveSubscriptionSet, 2098 auth: &BaseAuthState, 2099 projection: &RelayProjectionContext, 2100 ) -> Result<Vec<RuntimeRelayMessage>, BaseRelayError> { 2101 let pocket_event = self.inner.store.event_by_offset(offset.as_u64())?; 2102 let group_auth = GroupAuthContext::new(auth.authenticated_pubkeys().iter().cloned()); 2103 let mut messages = Vec::new(); 2104 let mut delivered = BTreeSet::new(); 2105 let mut delivery = RelayLiveProjectionDelivery { 2106 subscriptions, 2107 auth, 2108 projection, 2109 group_auth: &group_auth, 2110 delivered: &mut delivered, 2111 messages: &mut messages, 2112 }; 2113 self.fanout_projected_live_event(&pocket_event, offset.as_u64(), &mut delivery)?; 2114 let context = RelayLiveProjectionContext::new_with_authenticated_pubkeys( 2115 projection.clone(), 2116 offset.as_u64(), 2117 RelayEventContext::from_pocket_event(&pocket_event)?, 2118 auth.authenticated_pubkeys().iter().cloned().collect(), 2119 ); 2120 for candidate in self.inner.hooks.live_projection_candidates(&context) { 2121 let Ok(candidate_event) = self.inner.store.event_by_offset(candidate.store_offset()) 2122 else { 2123 continue; 2124 }; 2125 self.fanout_projected_live_event( 2126 &candidate_event, 2127 candidate.store_offset(), 2128 &mut delivery, 2129 )?; 2130 } 2131 Ok(messages) 2132 } 2133 2134 fn fanout_projected_live_event( 2135 &self, 2136 event: &PocketEvent, 2137 store_offset: u64, 2138 delivery: &mut RelayLiveProjectionDelivery<'_>, 2139 ) -> Result<(), BaseRelayError> { 2140 let subscriptions = 2141 delivery 2142 .subscriptions 2143 .fanout(event, delivery.group_auth, |event, auth| { 2144 BaseRelay::group_read_gate_visible_to_auth( 2145 self.inner.groups.as_ref(), 2146 event, 2147 auth, 2148 ) 2149 .unwrap_or(false) 2150 })?; 2151 for matched in subscriptions { 2152 let subscription_id = matched.subscription_id().clone(); 2153 if let Some(projected) = 2154 self.inner 2155 .project_event_output(RelayEventProjectionRequest { 2156 subscription_id: matched.subscription_id(), 2157 projection: delivery.projection, 2158 source: RelayEventProjectionSource::LiveFanout { store_offset }, 2159 event, 2160 auth: delivery.auth, 2161 matched_filters: matched 2162 .matched_filter_contexts() 2163 .into_iter() 2164 .zip(matched.filters()) 2165 .collect(), 2166 })? 2167 { 2168 let event_id = projected.id().as_hex_string(); 2169 if delivery 2170 .delivered 2171 .insert((subscription_id.clone(), event_id)) 2172 { 2173 delivery 2174 .messages 2175 .push(RuntimeRelayMessage::event(subscription_id, projected)); 2176 } 2177 } 2178 } 2179 Ok(()) 2180 } 2181 2182 pub async fn shutdown(&self) -> Result<BaseRelayShutdownReport, BaseRelayError> { 2183 self.inner.shutdown.request_shutdown(); 2184 self.inner.store.sync()?; 2185 Ok(BaseRelayShutdownReport::new(0)) 2186 } 2187 } 2188 2189 fn auth_response_failed(replies: &[RelayMessage]) -> bool { 2190 replies.iter().any(|reply| { 2191 matches!( 2192 reply, 2193 RelayMessage::Ok { 2194 accepted: false, 2195 .. 2196 } 2197 ) 2198 }) 2199 } 2200 2201 fn record_event_metrics( 2202 metrics: &TangleRuntimeMetrics, 2203 message: &RelayMessage, 2204 is_group_event: bool, 2205 started_at: Instant, 2206 ) { 2207 metrics.record_event_admission_latency(elapsed_micros(started_at)); 2208 if let RelayMessage::Ok { accepted, .. } = message { 2209 if *accepted { 2210 metrics.record_event_admission(); 2211 } else { 2212 metrics.record_event_rejection(); 2213 if is_group_event { 2214 metrics.record_group_write_denial(); 2215 } 2216 } 2217 } 2218 } 2219 2220 fn elapsed_micros(started_at: Instant) -> u64 { 2221 u64::try_from(started_at.elapsed().as_micros()).unwrap_or(u64::MAX) 2222 } 2223 2224 fn directory_size_bytes(path: &Path) -> u64 { 2225 let Ok(metadata) = fs::metadata(path) else { 2226 return 0; 2227 }; 2228 if metadata.is_file() { 2229 return metadata.len(); 2230 } 2231 if !metadata.is_dir() { 2232 return 0; 2233 } 2234 let Ok(entries) = fs::read_dir(path) else { 2235 return 0; 2236 }; 2237 entries 2238 .filter_map(Result::ok) 2239 .map(|entry| directory_size_bytes(&entry.path())) 2240 .sum() 2241 } 2242 2243 #[cfg(test)] 2244 fn protocol_client_message_to_runtime_for_test( 2245 message: tangle_protocol::ClientMessage, 2246 ) -> Result<RuntimeClientMessage, BaseRelayError> { 2247 match message { 2248 tangle_protocol::ClientMessage::Event(event) => Ok(RuntimeClientMessage::Event( 2249 crate::pocket_conversion::tangle_event_to_pocket(&event)?, 2250 )), 2251 tangle_protocol::ClientMessage::Req { 2252 subscription_id, 2253 filters, 2254 } => Ok(RuntimeClientMessage::Req { 2255 subscription_id, 2256 search_present: filters.iter().any(|filter| filter.search().is_some()), 2257 filters: filters 2258 .iter() 2259 .map(crate::pocket_conversion::tangle_filter_to_pocket) 2260 .collect::<Result<Vec<_>, _>>()?, 2261 }), 2262 tangle_protocol::ClientMessage::Count { 2263 subscription_id, 2264 filters, 2265 } => Ok(RuntimeClientMessage::Count { 2266 subscription_id, 2267 search_present: filters.iter().any(|filter| filter.search().is_some()), 2268 filters: filters 2269 .iter() 2270 .map(crate::pocket_conversion::tangle_filter_to_pocket) 2271 .collect::<Result<Vec<_>, _>>()?, 2272 }), 2273 tangle_protocol::ClientMessage::Close(subscription_id) => { 2274 Ok(RuntimeClientMessage::Close(subscription_id)) 2275 } 2276 tangle_protocol::ClientMessage::Auth(event) => Ok(RuntimeClientMessage::Auth( 2277 crate::pocket_conversion::tangle_event_to_pocket(&event)?, 2278 )), 2279 tangle_protocol::ClientMessage::NegOpen { 2280 subscription_id, 2281 filter, 2282 message, 2283 } => Ok(RuntimeClientMessage::NegOpen { 2284 subscription_id, 2285 filter: crate::pocket_conversion::tangle_filter_to_pocket(&filter)?, 2286 message, 2287 }), 2288 tangle_protocol::ClientMessage::NegMsg { 2289 subscription_id, 2290 message, 2291 } => Ok(RuntimeClientMessage::NegMsg { 2292 subscription_id, 2293 message, 2294 }), 2295 tangle_protocol::ClientMessage::NegClose(subscription_id) => { 2296 Ok(RuntimeClientMessage::NegClose(subscription_id)) 2297 } 2298 } 2299 } 2300 2301 fn runtime_client_message_metric_kind( 2302 message: &RuntimeClientMessage, 2303 ) -> TangleClientMessageMetricKind { 2304 match message { 2305 RuntimeClientMessage::Event(_) => TangleClientMessageMetricKind::Event, 2306 RuntimeClientMessage::Req { .. } => TangleClientMessageMetricKind::Req, 2307 RuntimeClientMessage::Count { .. } => TangleClientMessageMetricKind::Count, 2308 RuntimeClientMessage::Auth(_) => TangleClientMessageMetricKind::Auth, 2309 RuntimeClientMessage::Close(_) => TangleClientMessageMetricKind::Close, 2310 RuntimeClientMessage::NegOpen { .. } 2311 | RuntimeClientMessage::NegMsg { .. } 2312 | RuntimeClientMessage::NegClose(_) => TangleClientMessageMetricKind::Negentropy, 2313 } 2314 } 2315 2316 fn pocket_filter_group_ids(filters: &[PocketOwnedFilter]) -> Vec<GroupId> { 2317 let mut group_ids = BTreeSet::new(); 2318 for filter in filters { 2319 let Ok(tags) = filter.tags() else { 2320 continue; 2321 }; 2322 for mut tag in tags.iter() { 2323 let name = tag.next(); 2324 if !matches!(name, Some(value) if matches!(value, b"h" | b"d")) { 2325 continue; 2326 } 2327 for value in tag { 2328 if let Ok(value) = std::str::from_utf8(value) 2329 && let Ok(group_id) = GroupId::new(value) 2330 { 2331 group_ids.insert(group_id); 2332 } 2333 } 2334 } 2335 } 2336 group_ids.into_iter().collect() 2337 } 2338 2339 fn pocket_filter_kinds(filters: &[PocketOwnedFilter]) -> Vec<Kind> { 2340 filters 2341 .iter() 2342 .flat_map(|filter| filter.kinds()) 2343 .filter_map(|kind| Kind::new(u64::from(kind.as_u16())).ok()) 2344 .collect::<BTreeSet<_>>() 2345 .into_iter() 2346 .collect() 2347 } 2348 2349 impl fmt::Debug for RelayRuntimeHandle { 2350 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 2351 formatter.write_str("RelayRuntimeHandle") 2352 } 2353 } 2354 2355 #[derive(Debug, Clone, Copy, PartialEq, Eq)] 2356 pub struct TangleRuntimeLimits { 2357 max_message_length: usize, 2358 base_relay_limits: BaseRelayLimits, 2359 event_bus_capacity: usize, 2360 outbound_queue_capacity: usize, 2361 } 2362 2363 impl TangleRuntimeLimits { 2364 pub fn new( 2365 max_message_length: usize, 2366 base_relay_limits: BaseRelayLimits, 2367 event_bus_capacity: usize, 2368 outbound_queue_capacity: usize, 2369 ) -> Result<Self, BaseRelayError> { 2370 if max_message_length == 0 { 2371 return Err(BaseRelayError::invalid( 2372 "runtime max message length must be greater than zero", 2373 )); 2374 } 2375 if event_bus_capacity == 0 { 2376 return Err(BaseRelayError::invalid( 2377 "runtime event bus capacity must be greater than zero", 2378 )); 2379 } 2380 if outbound_queue_capacity == 0 { 2381 return Err(BaseRelayError::invalid( 2382 "runtime outbound queue capacity must be greater than zero", 2383 )); 2384 } 2385 Ok(Self { 2386 max_message_length, 2387 base_relay_limits, 2388 event_bus_capacity, 2389 outbound_queue_capacity, 2390 }) 2391 } 2392 2393 pub fn from_config(config: &BaseRelayRuntimeConfig) -> Result<Self, BaseRelayError> { 2394 let limits = config.limits(); 2395 Self::new( 2396 limits.max_message_length(), 2397 limits.base_relay_limits()?, 2398 limits.broadcast_channel_capacity(), 2399 limits.per_connection_outbound_queue(), 2400 ) 2401 } 2402 2403 pub fn max_message_length(self) -> usize { 2404 self.max_message_length 2405 } 2406 2407 pub fn base_relay_limits(self) -> BaseRelayLimits { 2408 self.base_relay_limits 2409 } 2410 2411 pub fn max_pending_events(self) -> usize { 2412 self.base_relay_limits.max_pending_events() 2413 } 2414 2415 pub fn event_bus_capacity(self) -> usize { 2416 self.event_bus_capacity 2417 } 2418 2419 pub fn outbound_queue_capacity(self) -> usize { 2420 self.outbound_queue_capacity 2421 } 2422 } 2423 2424 #[derive(Debug, Clone)] 2425 pub struct TangleRuntimeMetrics { 2426 inner: Arc<TangleRuntimeMetricsInner>, 2427 } 2428 2429 #[derive(Debug)] 2430 struct TangleRuntimeMetricsInner { 2431 started_at: Instant, 2432 active_sessions: AtomicUsize, 2433 total_sessions: AtomicU64, 2434 client_messages: AtomicU64, 2435 event_messages: AtomicU64, 2436 req_messages: AtomicU64, 2437 count_messages: AtomicU64, 2438 auth_messages: AtomicU64, 2439 close_messages: AtomicU64, 2440 active_subscriptions: AtomicUsize, 2441 opened_subscriptions: AtomicU64, 2442 closed_subscriptions: AtomicU64, 2443 stored_event_offsets: AtomicU64, 2444 rate_limit_rejections: AtomicU64, 2445 auth_successes: AtomicU64, 2446 auth_failures: AtomicU64, 2447 event_admissions: AtomicU64, 2448 event_rejections: AtomicU64, 2449 group_read_denials: AtomicU64, 2450 group_write_denials: AtomicU64, 2451 event_bus_receivers_current: AtomicUsize, 2452 event_bus_published_offsets: AtomicU64, 2453 event_bus_lagged_receivers: AtomicU64, 2454 event_bus_lagged_offsets: AtomicU64, 2455 outbound_queue_full_closes: AtomicU64, 2456 outbox_pending_events: AtomicUsize, 2457 outbox_replayed_events: AtomicU64, 2458 disk_used_bytes: AtomicU64, 2459 event_admission_latency_total_micros: AtomicU64, 2460 event_admission_latency_count: AtomicU64, 2461 query_latency_total_micros: AtomicU64, 2462 query_latency_count: AtomicU64, 2463 query_candidates_scanned: AtomicU64, 2464 query_returned_events: AtomicU64, 2465 query_redacted_events: AtomicU64, 2466 count_refusals: AtomicU64, 2467 broad_query_rejections: AtomicU64, 2468 } 2469 2470 #[derive(Debug, Clone, Copy, PartialEq, Eq)] 2471 pub enum TangleClientMessageMetricKind { 2472 Event, 2473 Req, 2474 Count, 2475 Auth, 2476 Close, 2477 Negentropy, 2478 } 2479 2480 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] 2481 pub struct TangleRuntimeMetricsSnapshot { 2482 tangle_runtime_uptime_seconds: u64, 2483 tangle_readiness_ready: bool, 2484 tangle_ws_connections_current: usize, 2485 tangle_ws_connections_total: u64, 2486 tangle_client_messages_total: u64, 2487 tangle_event_messages_total: u64, 2488 tangle_req_messages_total: u64, 2489 tangle_count_messages_total: u64, 2490 tangle_auth_messages_total: u64, 2491 tangle_close_messages_total: u64, 2492 tangle_subscriptions_current: usize, 2493 tangle_subscriptions_opened_total: u64, 2494 tangle_subscriptions_closed_total: u64, 2495 tangle_stored_event_offsets_total: u64, 2496 tangle_rate_limit_rejections_total: u64, 2497 tangle_auth_success_total: u64, 2498 tangle_auth_failure_total: u64, 2499 tangle_event_admitted_total: u64, 2500 tangle_event_rejected_total: u64, 2501 tangle_group_read_denied_total: u64, 2502 tangle_group_write_denied_total: u64, 2503 tangle_event_bus_receivers_current: usize, 2504 tangle_event_bus_published_offsets_total: u64, 2505 tangle_event_bus_lagged_receivers_total: u64, 2506 tangle_event_bus_lagged_offsets_total: u64, 2507 tangle_outbound_queue_full_closes_total: u64, 2508 tangle_outbox_pending_events: usize, 2509 tangle_outbox_replayed_events_total: u64, 2510 tangle_disk_used_bytes: u64, 2511 tangle_event_admission_latency_total_micros: u64, 2512 tangle_event_admission_latency_count: u64, 2513 tangle_query_latency_total_micros: u64, 2514 tangle_query_latency_count: u64, 2515 tangle_query_candidates_scanned_total: u64, 2516 tangle_query_returned_events_total: u64, 2517 tangle_query_redacted_events_total: u64, 2518 tangle_count_refusals_total: u64, 2519 tangle_broad_query_rejections_total: u64, 2520 } 2521 2522 impl TangleRuntimeMetricsSnapshot { 2523 pub fn prometheus_text(&self) -> String { 2524 let samples = [ 2525 ( 2526 "tangle_runtime_uptime_seconds", 2527 "gauge", 2528 self.tangle_runtime_uptime_seconds.to_string(), 2529 ), 2530 ( 2531 "tangle_readiness_ready", 2532 "gauge", 2533 u8::from(self.tangle_readiness_ready).to_string(), 2534 ), 2535 ( 2536 "tangle_ws_connections_current", 2537 "gauge", 2538 self.tangle_ws_connections_current.to_string(), 2539 ), 2540 ( 2541 "tangle_ws_connections_total", 2542 "counter", 2543 self.tangle_ws_connections_total.to_string(), 2544 ), 2545 ( 2546 "tangle_client_messages_total", 2547 "counter", 2548 self.tangle_client_messages_total.to_string(), 2549 ), 2550 ( 2551 "tangle_event_messages_total", 2552 "counter", 2553 self.tangle_event_messages_total.to_string(), 2554 ), 2555 ( 2556 "tangle_req_messages_total", 2557 "counter", 2558 self.tangle_req_messages_total.to_string(), 2559 ), 2560 ( 2561 "tangle_count_messages_total", 2562 "counter", 2563 self.tangle_count_messages_total.to_string(), 2564 ), 2565 ( 2566 "tangle_auth_messages_total", 2567 "counter", 2568 self.tangle_auth_messages_total.to_string(), 2569 ), 2570 ( 2571 "tangle_close_messages_total", 2572 "counter", 2573 self.tangle_close_messages_total.to_string(), 2574 ), 2575 ( 2576 "tangle_subscriptions_current", 2577 "gauge", 2578 self.tangle_subscriptions_current.to_string(), 2579 ), 2580 ( 2581 "tangle_subscriptions_opened_total", 2582 "counter", 2583 self.tangle_subscriptions_opened_total.to_string(), 2584 ), 2585 ( 2586 "tangle_subscriptions_closed_total", 2587 "counter", 2588 self.tangle_subscriptions_closed_total.to_string(), 2589 ), 2590 ( 2591 "tangle_stored_event_offsets_total", 2592 "counter", 2593 self.tangle_stored_event_offsets_total.to_string(), 2594 ), 2595 ( 2596 "tangle_rate_limit_rejections_total", 2597 "counter", 2598 self.tangle_rate_limit_rejections_total.to_string(), 2599 ), 2600 ( 2601 "tangle_auth_success_total", 2602 "counter", 2603 self.tangle_auth_success_total.to_string(), 2604 ), 2605 ( 2606 "tangle_auth_failure_total", 2607 "counter", 2608 self.tangle_auth_failure_total.to_string(), 2609 ), 2610 ( 2611 "tangle_event_admitted_total", 2612 "counter", 2613 self.tangle_event_admitted_total.to_string(), 2614 ), 2615 ( 2616 "tangle_event_rejected_total", 2617 "counter", 2618 self.tangle_event_rejected_total.to_string(), 2619 ), 2620 ( 2621 "tangle_group_read_denied_total", 2622 "counter", 2623 self.tangle_group_read_denied_total.to_string(), 2624 ), 2625 ( 2626 "tangle_group_write_denied_total", 2627 "counter", 2628 self.tangle_group_write_denied_total.to_string(), 2629 ), 2630 ( 2631 "tangle_event_bus_receivers_current", 2632 "gauge", 2633 self.tangle_event_bus_receivers_current.to_string(), 2634 ), 2635 ( 2636 "tangle_event_bus_published_offsets_total", 2637 "counter", 2638 self.tangle_event_bus_published_offsets_total.to_string(), 2639 ), 2640 ( 2641 "tangle_event_bus_lagged_receivers_total", 2642 "counter", 2643 self.tangle_event_bus_lagged_receivers_total.to_string(), 2644 ), 2645 ( 2646 "tangle_event_bus_lagged_offsets_total", 2647 "counter", 2648 self.tangle_event_bus_lagged_offsets_total.to_string(), 2649 ), 2650 ( 2651 "tangle_outbound_queue_full_closes_total", 2652 "counter", 2653 self.tangle_outbound_queue_full_closes_total.to_string(), 2654 ), 2655 ( 2656 "tangle_outbox_pending_events", 2657 "gauge", 2658 self.tangle_outbox_pending_events.to_string(), 2659 ), 2660 ( 2661 "tangle_outbox_replayed_events_total", 2662 "counter", 2663 self.tangle_outbox_replayed_events_total.to_string(), 2664 ), 2665 ( 2666 "tangle_disk_used_bytes", 2667 "gauge", 2668 self.tangle_disk_used_bytes.to_string(), 2669 ), 2670 ( 2671 "tangle_event_admission_latency_total_micros", 2672 "counter", 2673 self.tangle_event_admission_latency_total_micros.to_string(), 2674 ), 2675 ( 2676 "tangle_event_admission_latency_count", 2677 "counter", 2678 self.tangle_event_admission_latency_count.to_string(), 2679 ), 2680 ( 2681 "tangle_query_latency_total_micros", 2682 "counter", 2683 self.tangle_query_latency_total_micros.to_string(), 2684 ), 2685 ( 2686 "tangle_query_latency_count", 2687 "counter", 2688 self.tangle_query_latency_count.to_string(), 2689 ), 2690 ( 2691 "tangle_query_candidates_scanned_total", 2692 "counter", 2693 self.tangle_query_candidates_scanned_total.to_string(), 2694 ), 2695 ( 2696 "tangle_query_returned_events_total", 2697 "counter", 2698 self.tangle_query_returned_events_total.to_string(), 2699 ), 2700 ( 2701 "tangle_query_redacted_events_total", 2702 "counter", 2703 self.tangle_query_redacted_events_total.to_string(), 2704 ), 2705 ( 2706 "tangle_count_refusals_total", 2707 "counter", 2708 self.tangle_count_refusals_total.to_string(), 2709 ), 2710 ( 2711 "tangle_broad_query_rejections_total", 2712 "counter", 2713 self.tangle_broad_query_rejections_total.to_string(), 2714 ), 2715 ]; 2716 let mut output = String::with_capacity(samples.len() * 128); 2717 for (name, metric_type, value) in samples { 2718 output.push_str("# TYPE "); 2719 output.push_str(name); 2720 output.push(' '); 2721 output.push_str(metric_type); 2722 output.push('\n'); 2723 output.push_str(name); 2724 output.push(' '); 2725 output.push_str(&value); 2726 output.push('\n'); 2727 } 2728 output 2729 } 2730 2731 pub fn active_sessions(&self) -> usize { 2732 self.tangle_ws_connections_current 2733 } 2734 2735 pub fn total_sessions(&self) -> u64 { 2736 self.tangle_ws_connections_total 2737 } 2738 2739 pub fn client_messages(&self) -> u64 { 2740 self.tangle_client_messages_total 2741 } 2742 2743 pub fn event_messages(&self) -> u64 { 2744 self.tangle_event_messages_total 2745 } 2746 2747 pub fn req_messages(&self) -> u64 { 2748 self.tangle_req_messages_total 2749 } 2750 2751 pub fn count_messages(&self) -> u64 { 2752 self.tangle_count_messages_total 2753 } 2754 2755 pub fn auth_messages(&self) -> u64 { 2756 self.tangle_auth_messages_total 2757 } 2758 2759 pub fn close_messages(&self) -> u64 { 2760 self.tangle_close_messages_total 2761 } 2762 2763 pub fn opened_subscriptions(&self) -> u64 { 2764 self.tangle_subscriptions_opened_total 2765 } 2766 2767 pub fn active_subscriptions(&self) -> usize { 2768 self.tangle_subscriptions_current 2769 } 2770 2771 pub fn closed_subscriptions(&self) -> u64 { 2772 self.tangle_subscriptions_closed_total 2773 } 2774 2775 pub fn stored_event_offsets(&self) -> u64 { 2776 self.tangle_stored_event_offsets_total 2777 } 2778 2779 pub fn rate_limit_rejections(&self) -> u64 { 2780 self.tangle_rate_limit_rejections_total 2781 } 2782 } 2783 2784 impl TangleRuntimeMetrics { 2785 pub fn new() -> Self { 2786 Self { 2787 inner: Arc::new(TangleRuntimeMetricsInner { 2788 started_at: Instant::now(), 2789 active_sessions: AtomicUsize::new(0), 2790 total_sessions: AtomicU64::new(0), 2791 client_messages: AtomicU64::new(0), 2792 event_messages: AtomicU64::new(0), 2793 req_messages: AtomicU64::new(0), 2794 count_messages: AtomicU64::new(0), 2795 auth_messages: AtomicU64::new(0), 2796 close_messages: AtomicU64::new(0), 2797 active_subscriptions: AtomicUsize::new(0), 2798 opened_subscriptions: AtomicU64::new(0), 2799 closed_subscriptions: AtomicU64::new(0), 2800 stored_event_offsets: AtomicU64::new(0), 2801 rate_limit_rejections: AtomicU64::new(0), 2802 auth_successes: AtomicU64::new(0), 2803 auth_failures: AtomicU64::new(0), 2804 event_admissions: AtomicU64::new(0), 2805 event_rejections: AtomicU64::new(0), 2806 group_read_denials: AtomicU64::new(0), 2807 group_write_denials: AtomicU64::new(0), 2808 event_bus_receivers_current: AtomicUsize::new(0), 2809 event_bus_published_offsets: AtomicU64::new(0), 2810 event_bus_lagged_receivers: AtomicU64::new(0), 2811 event_bus_lagged_offsets: AtomicU64::new(0), 2812 outbound_queue_full_closes: AtomicU64::new(0), 2813 outbox_pending_events: AtomicUsize::new(0), 2814 outbox_replayed_events: AtomicU64::new(0), 2815 disk_used_bytes: AtomicU64::new(0), 2816 event_admission_latency_total_micros: AtomicU64::new(0), 2817 event_admission_latency_count: AtomicU64::new(0), 2818 query_latency_total_micros: AtomicU64::new(0), 2819 query_latency_count: AtomicU64::new(0), 2820 query_candidates_scanned: AtomicU64::new(0), 2821 query_returned_events: AtomicU64::new(0), 2822 query_redacted_events: AtomicU64::new(0), 2823 count_refusals: AtomicU64::new(0), 2824 broad_query_rejections: AtomicU64::new(0), 2825 }), 2826 } 2827 } 2828 2829 pub fn snapshot(&self) -> TangleRuntimeMetricsSnapshot { 2830 self.snapshot_with_readiness(false) 2831 } 2832 2833 pub fn snapshot_with_readiness(&self, readiness_ready: bool) -> TangleRuntimeMetricsSnapshot { 2834 TangleRuntimeMetricsSnapshot { 2835 tangle_runtime_uptime_seconds: self.started_at().elapsed().as_secs(), 2836 tangle_readiness_ready: readiness_ready, 2837 tangle_ws_connections_current: self.active_sessions(), 2838 tangle_ws_connections_total: self.total_sessions(), 2839 tangle_client_messages_total: self.client_messages(), 2840 tangle_event_messages_total: self.event_messages(), 2841 tangle_req_messages_total: self.req_messages(), 2842 tangle_count_messages_total: self.count_messages(), 2843 tangle_auth_messages_total: self.auth_messages(), 2844 tangle_close_messages_total: self.close_messages(), 2845 tangle_subscriptions_current: self.active_subscriptions(), 2846 tangle_subscriptions_opened_total: self.opened_subscriptions(), 2847 tangle_subscriptions_closed_total: self.closed_subscriptions(), 2848 tangle_stored_event_offsets_total: self.stored_event_offsets(), 2849 tangle_rate_limit_rejections_total: self.rate_limit_rejections(), 2850 tangle_auth_success_total: self.auth_successes(), 2851 tangle_auth_failure_total: self.auth_failures(), 2852 tangle_event_admitted_total: self.event_admissions(), 2853 tangle_event_rejected_total: self.event_rejections(), 2854 tangle_group_read_denied_total: self.group_read_denials(), 2855 tangle_group_write_denied_total: self.group_write_denials(), 2856 tangle_event_bus_receivers_current: self.event_bus_receivers_current(), 2857 tangle_event_bus_published_offsets_total: self.event_bus_published_offsets(), 2858 tangle_event_bus_lagged_receivers_total: self.event_bus_lagged_receivers(), 2859 tangle_event_bus_lagged_offsets_total: self.event_bus_lagged_offsets(), 2860 tangle_outbound_queue_full_closes_total: self.outbound_queue_full_closes(), 2861 tangle_outbox_pending_events: self.outbox_pending_events(), 2862 tangle_outbox_replayed_events_total: self.outbox_replayed_events(), 2863 tangle_disk_used_bytes: self.disk_used_bytes(), 2864 tangle_event_admission_latency_total_micros: self 2865 .event_admission_latency_total_micros(), 2866 tangle_event_admission_latency_count: self.event_admission_latency_count(), 2867 tangle_query_latency_total_micros: self.query_latency_total_micros(), 2868 tangle_query_latency_count: self.query_latency_count(), 2869 tangle_query_candidates_scanned_total: self.query_candidates_scanned(), 2870 tangle_query_returned_events_total: self.query_returned_events(), 2871 tangle_query_redacted_events_total: self.query_redacted_events(), 2872 tangle_count_refusals_total: self.count_refusals(), 2873 tangle_broad_query_rejections_total: self.broad_query_rejections(), 2874 } 2875 } 2876 2877 pub fn started_at(&self) -> Instant { 2878 self.inner.started_at 2879 } 2880 2881 pub fn active_sessions(&self) -> usize { 2882 self.inner.active_sessions.load(Ordering::Relaxed) 2883 } 2884 2885 pub fn total_sessions(&self) -> u64 { 2886 self.inner.total_sessions.load(Ordering::Relaxed) 2887 } 2888 2889 pub fn client_messages(&self) -> u64 { 2890 self.inner.client_messages.load(Ordering::Relaxed) 2891 } 2892 2893 pub fn event_messages(&self) -> u64 { 2894 self.inner.event_messages.load(Ordering::Relaxed) 2895 } 2896 2897 pub fn req_messages(&self) -> u64 { 2898 self.inner.req_messages.load(Ordering::Relaxed) 2899 } 2900 2901 pub fn count_messages(&self) -> u64 { 2902 self.inner.count_messages.load(Ordering::Relaxed) 2903 } 2904 2905 pub fn auth_messages(&self) -> u64 { 2906 self.inner.auth_messages.load(Ordering::Relaxed) 2907 } 2908 2909 pub fn close_messages(&self) -> u64 { 2910 self.inner.close_messages.load(Ordering::Relaxed) 2911 } 2912 2913 pub fn opened_subscriptions(&self) -> u64 { 2914 self.inner.opened_subscriptions.load(Ordering::Relaxed) 2915 } 2916 2917 pub fn active_subscriptions(&self) -> usize { 2918 self.inner.active_subscriptions.load(Ordering::Relaxed) 2919 } 2920 2921 pub fn closed_subscriptions(&self) -> u64 { 2922 self.inner.closed_subscriptions.load(Ordering::Relaxed) 2923 } 2924 2925 pub fn stored_event_offsets(&self) -> u64 { 2926 self.inner.stored_event_offsets.load(Ordering::Relaxed) 2927 } 2928 2929 pub fn rate_limit_rejections(&self) -> u64 { 2930 self.inner.rate_limit_rejections.load(Ordering::Relaxed) 2931 } 2932 2933 pub fn auth_successes(&self) -> u64 { 2934 self.inner.auth_successes.load(Ordering::Relaxed) 2935 } 2936 2937 pub fn auth_failures(&self) -> u64 { 2938 self.inner.auth_failures.load(Ordering::Relaxed) 2939 } 2940 2941 pub fn event_admissions(&self) -> u64 { 2942 self.inner.event_admissions.load(Ordering::Relaxed) 2943 } 2944 2945 pub fn event_rejections(&self) -> u64 { 2946 self.inner.event_rejections.load(Ordering::Relaxed) 2947 } 2948 2949 pub fn group_read_denials(&self) -> u64 { 2950 self.inner.group_read_denials.load(Ordering::Relaxed) 2951 } 2952 2953 pub fn group_write_denials(&self) -> u64 { 2954 self.inner.group_write_denials.load(Ordering::Relaxed) 2955 } 2956 2957 pub fn event_bus_receivers_current(&self) -> usize { 2958 self.inner 2959 .event_bus_receivers_current 2960 .load(Ordering::Relaxed) 2961 } 2962 2963 pub fn event_bus_published_offsets(&self) -> u64 { 2964 self.inner 2965 .event_bus_published_offsets 2966 .load(Ordering::Relaxed) 2967 } 2968 2969 pub fn event_bus_lagged_receivers(&self) -> u64 { 2970 self.inner 2971 .event_bus_lagged_receivers 2972 .load(Ordering::Relaxed) 2973 } 2974 2975 pub fn event_bus_lagged_offsets(&self) -> u64 { 2976 self.inner.event_bus_lagged_offsets.load(Ordering::Relaxed) 2977 } 2978 2979 pub fn outbound_queue_full_closes(&self) -> u64 { 2980 self.inner 2981 .outbound_queue_full_closes 2982 .load(Ordering::Relaxed) 2983 } 2984 2985 pub fn outbox_pending_events(&self) -> usize { 2986 self.inner.outbox_pending_events.load(Ordering::Relaxed) 2987 } 2988 2989 pub fn outbox_replayed_events(&self) -> u64 { 2990 self.inner.outbox_replayed_events.load(Ordering::Relaxed) 2991 } 2992 2993 pub fn disk_used_bytes(&self) -> u64 { 2994 self.inner.disk_used_bytes.load(Ordering::Relaxed) 2995 } 2996 2997 pub fn event_admission_latency_total_micros(&self) -> u64 { 2998 self.inner 2999 .event_admission_latency_total_micros 3000 .load(Ordering::Relaxed) 3001 } 3002 3003 pub fn event_admission_latency_count(&self) -> u64 { 3004 self.inner 3005 .event_admission_latency_count 3006 .load(Ordering::Relaxed) 3007 } 3008 3009 pub fn query_latency_total_micros(&self) -> u64 { 3010 self.inner 3011 .query_latency_total_micros 3012 .load(Ordering::Relaxed) 3013 } 3014 3015 pub fn query_latency_count(&self) -> u64 { 3016 self.inner.query_latency_count.load(Ordering::Relaxed) 3017 } 3018 3019 pub fn query_candidates_scanned(&self) -> u64 { 3020 self.inner.query_candidates_scanned.load(Ordering::Relaxed) 3021 } 3022 3023 pub fn query_returned_events(&self) -> u64 { 3024 self.inner.query_returned_events.load(Ordering::Relaxed) 3025 } 3026 3027 pub fn query_redacted_events(&self) -> u64 { 3028 self.inner.query_redacted_events.load(Ordering::Relaxed) 3029 } 3030 3031 pub fn count_refusals(&self) -> u64 { 3032 self.inner.count_refusals.load(Ordering::Relaxed) 3033 } 3034 3035 pub fn broad_query_rejections(&self) -> u64 { 3036 self.inner.broad_query_rejections.load(Ordering::Relaxed) 3037 } 3038 3039 pub fn record_session_opened(&self) -> usize { 3040 self.inner.total_sessions.fetch_add(1, Ordering::Relaxed); 3041 self.inner.active_sessions.fetch_add(1, Ordering::Relaxed) + 1 3042 } 3043 3044 pub fn record_session_closed(&self) -> usize { 3045 let mut current = self.inner.active_sessions.load(Ordering::Relaxed); 3046 loop { 3047 if current == 0 { 3048 return 0; 3049 } 3050 match self.inner.active_sessions.compare_exchange( 3051 current, 3052 current - 1, 3053 Ordering::Relaxed, 3054 Ordering::Relaxed, 3055 ) { 3056 Ok(_) => return current - 1, 3057 Err(actual) => current = actual, 3058 } 3059 } 3060 } 3061 3062 pub fn record_client_message(&self, kind: TangleClientMessageMetricKind) -> u64 { 3063 let total = self.inner.client_messages.fetch_add(1, Ordering::Relaxed) + 1; 3064 match kind { 3065 TangleClientMessageMetricKind::Event => { 3066 self.inner.event_messages.fetch_add(1, Ordering::Relaxed); 3067 } 3068 TangleClientMessageMetricKind::Req => { 3069 self.inner.req_messages.fetch_add(1, Ordering::Relaxed); 3070 } 3071 TangleClientMessageMetricKind::Count => { 3072 self.inner.count_messages.fetch_add(1, Ordering::Relaxed); 3073 } 3074 TangleClientMessageMetricKind::Auth => { 3075 self.inner.auth_messages.fetch_add(1, Ordering::Relaxed); 3076 } 3077 TangleClientMessageMetricKind::Close => { 3078 self.inner.close_messages.fetch_add(1, Ordering::Relaxed); 3079 } 3080 TangleClientMessageMetricKind::Negentropy => {} 3081 }; 3082 total 3083 } 3084 3085 pub fn record_subscription_opened(&self) -> u64 { 3086 self.inner 3087 .active_subscriptions 3088 .fetch_add(1, Ordering::Relaxed); 3089 self.inner 3090 .opened_subscriptions 3091 .fetch_add(1, Ordering::Relaxed) 3092 + 1 3093 } 3094 3095 pub fn record_subscriptions_closed(&self, count: usize) -> u64 { 3096 let mut current = self.active_subscriptions(); 3097 loop { 3098 let remaining = current.saturating_sub(count); 3099 match self.inner.active_subscriptions.compare_exchange( 3100 current, 3101 remaining, 3102 Ordering::Relaxed, 3103 Ordering::Relaxed, 3104 ) { 3105 Ok(_) => break, 3106 Err(actual) => current = actual, 3107 } 3108 } 3109 self.inner.closed_subscriptions.fetch_add( 3110 u64::try_from(count).expect("subscription count fits in u64"), 3111 Ordering::Relaxed, 3112 ) + u64::try_from(count).expect("subscription count fits in u64") 3113 } 3114 3115 pub fn record_stored_event_offset(&self) -> u64 { 3116 self.inner 3117 .stored_event_offsets 3118 .fetch_add(1, Ordering::Relaxed) 3119 + 1 3120 } 3121 3122 pub fn record_rate_limit_rejection(&self) -> u64 { 3123 self.inner 3124 .rate_limit_rejections 3125 .fetch_add(1, Ordering::Relaxed) 3126 + 1 3127 } 3128 3129 pub fn record_auth_success(&self) -> u64 { 3130 self.inner.auth_successes.fetch_add(1, Ordering::Relaxed) + 1 3131 } 3132 3133 pub fn record_auth_failure(&self) -> u64 { 3134 self.inner.auth_failures.fetch_add(1, Ordering::Relaxed) + 1 3135 } 3136 3137 pub fn record_event_admission(&self) -> u64 { 3138 self.inner.event_admissions.fetch_add(1, Ordering::Relaxed) + 1 3139 } 3140 3141 pub fn record_event_rejection(&self) -> u64 { 3142 self.inner.event_rejections.fetch_add(1, Ordering::Relaxed) + 1 3143 } 3144 3145 pub fn record_group_read_denial(&self) -> u64 { 3146 self.inner 3147 .group_read_denials 3148 .fetch_add(1, Ordering::Relaxed) 3149 + 1 3150 } 3151 3152 pub fn record_group_write_denial(&self) -> u64 { 3153 self.inner 3154 .group_write_denials 3155 .fetch_add(1, Ordering::Relaxed) 3156 + 1 3157 } 3158 3159 pub fn record_event_bus_receivers(&self, count: usize) { 3160 self.inner 3161 .event_bus_receivers_current 3162 .store(count, Ordering::Relaxed); 3163 } 3164 3165 pub fn record_event_bus_publish(&self, receivers: usize) -> u64 { 3166 self.record_event_bus_receivers(receivers); 3167 self.inner 3168 .event_bus_published_offsets 3169 .fetch_add(1, Ordering::Relaxed) 3170 + 1 3171 } 3172 3173 pub fn record_event_bus_lagged(&self, skipped: u64) { 3174 self.inner 3175 .event_bus_lagged_receivers 3176 .fetch_add(1, Ordering::Relaxed); 3177 self.inner 3178 .event_bus_lagged_offsets 3179 .fetch_add(skipped, Ordering::Relaxed); 3180 } 3181 3182 pub fn record_outbound_queue_full_close(&self) -> u64 { 3183 self.inner 3184 .outbound_queue_full_closes 3185 .fetch_add(1, Ordering::Relaxed) 3186 + 1 3187 } 3188 3189 pub fn record_outbox_pending_events(&self, count: usize) { 3190 self.inner 3191 .outbox_pending_events 3192 .store(count, Ordering::Relaxed); 3193 } 3194 3195 pub fn record_outbox_replayed_event(&self) -> u64 { 3196 self.inner 3197 .outbox_replayed_events 3198 .fetch_add(1, Ordering::Relaxed) 3199 + 1 3200 } 3201 3202 pub fn record_disk_used_bytes(&self, bytes: u64) { 3203 self.inner.disk_used_bytes.store(bytes, Ordering::Relaxed); 3204 } 3205 3206 pub fn record_event_admission_latency(&self, micros: u64) { 3207 self.inner 3208 .event_admission_latency_total_micros 3209 .fetch_add(micros, Ordering::Relaxed); 3210 self.inner 3211 .event_admission_latency_count 3212 .fetch_add(1, Ordering::Relaxed); 3213 } 3214 3215 pub fn record_query_latency(&self, micros: u64) { 3216 self.inner 3217 .query_latency_total_micros 3218 .fetch_add(micros, Ordering::Relaxed); 3219 self.inner 3220 .query_latency_count 3221 .fetch_add(1, Ordering::Relaxed); 3222 } 3223 3224 pub(crate) fn record_query_metrics(&self, metrics: BaseRelayQueryMetrics) { 3225 self.inner 3226 .query_candidates_scanned 3227 .fetch_add(metrics.candidates_scanned(), Ordering::Relaxed); 3228 self.inner 3229 .query_returned_events 3230 .fetch_add(metrics.returned_events(), Ordering::Relaxed); 3231 self.inner 3232 .query_redacted_events 3233 .fetch_add(metrics.redacted_events(), Ordering::Relaxed); 3234 } 3235 3236 pub fn record_count_refusal(&self) -> u64 { 3237 self.inner.count_refusals.fetch_add(1, Ordering::Relaxed) + 1 3238 } 3239 3240 pub fn record_broad_query_rejection(&self) -> u64 { 3241 self.inner 3242 .broad_query_rejections 3243 .fetch_add(1, Ordering::Relaxed) 3244 + 1 3245 } 3246 } 3247 3248 impl Default for TangleRuntimeMetrics { 3249 fn default() -> Self { 3250 Self::new() 3251 } 3252 } 3253 3254 #[derive(Debug, Clone)] 3255 pub struct TangleShutdownSignal { 3256 sender: watch::Sender<bool>, 3257 } 3258 3259 impl TangleShutdownSignal { 3260 pub fn new() -> Self { 3261 let (sender, _) = watch::channel(false); 3262 Self { sender } 3263 } 3264 3265 pub fn subscribe(&self) -> watch::Receiver<bool> { 3266 self.sender.subscribe() 3267 } 3268 3269 pub fn request_shutdown(&self) { 3270 self.sender.send_replace(true); 3271 } 3272 3273 pub fn requested(&self) -> bool { 3274 *self.sender.borrow() 3275 } 3276 } 3277 3278 impl Default for TangleShutdownSignal { 3279 fn default() -> Self { 3280 Self::new() 3281 } 3282 } 3283 3284 #[cfg(test)] 3285 mod tests { 3286 use super::{ 3287 BROAD_QUERY_TIME_WINDOW_SECONDS, EventAdmissionDecision, RelayEventAdmissionContext, 3288 RelayEventProjectionContext, RelayEventProjectionDecision, RelayEventProjectionSource, 3289 RelayEventStoredContext, RelayLiveProjectionCandidate, RelayLiveProjectionContext, 3290 RelayProjectionContext, RelayProjectionQueryPlan, RelayQueryProjectionContext, 3291 RelayRequestedKinds, RelayRuntime, RelayRuntimeHandle, RelayRuntimeHooks, 3292 RuntimeClientMessage, TangleBroadQueryReason, TangleClientRateLimitContext, 3293 TangleQueryClassification, TangleQueryClassifier, TangleRuntimeLimits, 3294 }; 3295 use crate::config::{BaseRelayRuntimeConfig, parse_base_relay_runtime_config_json}; 3296 use crate::event_bus::{TangleEventBus, TangleEventReceiveError, TangleEventReceiver}; 3297 use crate::rate_limits::{TangleRateLimitKey, TangleRateLimitQueryClass, TangleRateLimitScope}; 3298 use crate::relay::auth::BaseAuthState; 3299 use crate::relay::core::{BaseRelayLimitSettings, BaseRelayLimits, BaseRelayQueryMetrics}; 3300 use crate::relay::live::LiveSubscriptionSet; 3301 use crate::relay::outbound::RuntimeRelayMessage; 3302 use serde_json::json; 3303 use std::{ 3304 collections::{BTreeMap, BTreeSet}, 3305 net::{IpAddr, Ipv4Addr}, 3306 num::NonZeroU32, 3307 path::{Path, PathBuf}, 3308 sync::{Arc, Mutex, mpsc}, 3309 time::Duration, 3310 }; 3311 use tangle_crypto::RelaySigner; 3312 use tangle_groups::{ 3313 CanonicalGroupEvent, GroupEventClass, GroupId, GroupProjection, KIND_GROUP_ADMINS, 3314 KIND_GROUP_CREATE_GROUP, KIND_GROUP_DELETE_GROUP, KIND_GROUP_JOIN_REQUEST, 3315 KIND_GROUP_LEAVE_REQUEST, KIND_GROUP_MEMBERS, KIND_GROUP_METADATA, KIND_GROUP_PUT_USER, 3316 KIND_GROUP_REMOVE_USER, MemberStatus, StoreOffset, rebuild_group_projection, 3317 }; 3318 use tangle_protocol::{ 3319 ClientMessage, Event, EventId, Filter, Kind, PublicKeyHex, RelayMessage, SignatureHex, 3320 SubscriptionId, Tag, UnixTimestamp, UnsignedEvent, filter_from_value, 3321 }; 3322 use tangle_store_pocket::{ 3323 PocketEvent, PocketKind, PocketOwnedEvent, PocketOwnedTags, PocketTime, 3324 }; 3325 use tangle_test_support::FixtureKey; 3326 3327 #[test] 3328 fn tangle_runtime_opens_owned_process_shell_from_config() { 3329 let root = temp_root("owned-runtime"); 3330 let _ = std::fs::remove_dir_all(&root); 3331 let config = runtime_config(&root, 8); 3332 3333 let mut runtime = RelayRuntime::open(config).expect("runtime"); 3334 let mut offsets = runtime.event_bus().subscribe(); 3335 let shutdown = runtime.shutdown_signal().subscribe(); 3336 3337 assert_eq!(runtime.config().relay_url(), "wss://relay.radroots.test"); 3338 assert_eq!(runtime.config().listen_addr().to_string(), "127.0.0.1:0"); 3339 assert_eq!(runtime.limits().max_pending_events(), 8); 3340 assert_eq!(runtime.limits().max_message_length(), 1_048_576); 3341 assert_eq!(runtime.limits().event_bus_capacity(), 16); 3342 assert_eq!(runtime.limits().outbound_queue_capacity(), 8); 3343 assert_eq!(runtime.event_bus().capacity(), 16); 3344 assert_eq!(runtime.event_bus().receiver_count(), 1); 3345 assert_eq!(runtime.rate_limiter().tracked_key_count(), 0); 3346 assert_eq!(runtime.metrics().active_sessions(), 0); 3347 assert_eq!(runtime.metrics().stored_event_offsets(), 0); 3348 assert!(runtime.relay().groups_enabled()); 3349 assert!(!runtime.readiness_state().is_ready()); 3350 assert_eq!( 3351 runtime.readiness_state().response().checks.server_bind, 3352 "not_ready" 3353 ); 3354 assert_eq!( 3355 runtime 3356 .readiness_state() 3357 .response() 3358 .checks 3359 .group_outbox_replay, 3360 "ready" 3361 ); 3362 assert_eq!( 3363 runtime.readiness_state().response().checks.event_bus, 3364 "ready" 3365 ); 3366 assert!(!*shutdown.borrow()); 3367 3368 assert_eq!(runtime.event_bus().publish(StoreOffset::new(42)), 1); 3369 assert_eq!(offsets.try_recv().expect("offset"), StoreOffset::new(42)); 3370 assert_eq!( 3371 runtime 3372 .auth_state() 3373 .expect("auth") 3374 .authenticated_pubkeys() 3375 .len(), 3376 0 3377 ); 3378 assert_eq!( 3379 runtime.config().pocket_config().data_directory(), 3380 Path::new(&root).join("pocket") 3381 ); 3382 3383 assert_eq!(runtime.metrics().record_session_opened(), 1); 3384 assert_eq!(runtime.metrics().active_sessions(), 1); 3385 assert_eq!(runtime.metrics().total_sessions(), 1); 3386 assert_eq!(runtime.metrics().record_session_closed(), 0); 3387 assert_eq!(runtime.metrics().active_sessions(), 0); 3388 assert_eq!(runtime.metrics().total_sessions(), 1); 3389 assert_eq!( 3390 runtime 3391 .metrics() 3392 .record_client_message(super::TangleClientMessageMetricKind::Req), 3393 1 3394 ); 3395 assert_eq!(runtime.metrics().client_messages(), 1); 3396 assert_eq!(runtime.metrics().req_messages(), 1); 3397 assert_eq!(runtime.metrics().record_subscription_opened(), 1); 3398 assert_eq!(runtime.metrics().active_subscriptions(), 1); 3399 assert_eq!(runtime.metrics().opened_subscriptions(), 1); 3400 assert_eq!(runtime.metrics().record_subscriptions_closed(1), 1); 3401 assert_eq!(runtime.metrics().active_subscriptions(), 0); 3402 assert_eq!(runtime.metrics().closed_subscriptions(), 1); 3403 assert_eq!(runtime.metrics().record_stored_event_offset(), 1); 3404 assert_eq!(runtime.metrics().stored_event_offsets(), 1); 3405 assert_eq!(runtime.metrics().record_rate_limit_rejection(), 1); 3406 assert_eq!(runtime.metrics().rate_limit_rejections(), 1); 3407 assert_eq!(runtime.metrics().record_auth_success(), 1); 3408 assert_eq!(runtime.metrics().record_auth_failure(), 1); 3409 assert_eq!(runtime.metrics().record_event_admission(), 1); 3410 assert_eq!(runtime.metrics().record_event_rejection(), 1); 3411 assert_eq!(runtime.metrics().record_group_read_denial(), 1); 3412 assert_eq!(runtime.metrics().record_group_write_denial(), 1); 3413 runtime.metrics().record_event_bus_receivers(3); 3414 assert_eq!(runtime.metrics().record_event_bus_publish(3), 1); 3415 runtime.metrics().record_event_bus_lagged(4); 3416 assert_eq!(runtime.metrics().record_outbound_queue_full_close(), 1); 3417 runtime.metrics().record_outbox_pending_events(2); 3418 assert_eq!(runtime.metrics().record_outbox_replayed_event(), 1); 3419 runtime.metrics().record_disk_used_bytes(5); 3420 runtime.metrics().record_event_admission_latency(13); 3421 runtime.metrics().record_query_latency(17); 3422 runtime 3423 .metrics() 3424 .record_query_metrics(BaseRelayQueryMetrics::new(5, 3, 2)); 3425 assert_eq!(runtime.metrics().record_count_refusal(), 1); 3426 assert_eq!(runtime.metrics().record_broad_query_rejection(), 1); 3427 let snapshot = runtime.metrics().snapshot_with_readiness(true); 3428 assert_eq!(snapshot.active_sessions(), 0); 3429 assert_eq!(snapshot.total_sessions(), 1); 3430 assert_eq!(snapshot.client_messages(), 1); 3431 assert_eq!(snapshot.req_messages(), 1); 3432 assert_eq!(snapshot.active_subscriptions(), 0); 3433 assert_eq!(snapshot.opened_subscriptions(), 1); 3434 assert_eq!(snapshot.closed_subscriptions(), 1); 3435 assert_eq!(snapshot.stored_event_offsets(), 1); 3436 assert_eq!(snapshot.rate_limit_rejections(), 1); 3437 let snapshot_value = serde_json::to_value(snapshot).expect("snapshot json"); 3438 assert_eq!(snapshot_value["tangle_readiness_ready"], true); 3439 assert_eq!(snapshot_value["tangle_auth_success_total"], 1); 3440 assert_eq!(snapshot_value["tangle_auth_failure_total"], 1); 3441 assert_eq!(snapshot_value["tangle_event_admitted_total"], 1); 3442 assert_eq!(snapshot_value["tangle_event_rejected_total"], 1); 3443 assert_eq!(snapshot_value["tangle_group_read_denied_total"], 1); 3444 assert_eq!(snapshot_value["tangle_group_write_denied_total"], 1); 3445 assert_eq!(snapshot_value["tangle_event_bus_receivers_current"], 3); 3446 assert_eq!( 3447 snapshot_value["tangle_event_bus_published_offsets_total"], 3448 1 3449 ); 3450 assert_eq!(snapshot_value["tangle_event_bus_lagged_receivers_total"], 1); 3451 assert_eq!(snapshot_value["tangle_event_bus_lagged_offsets_total"], 4); 3452 assert_eq!(snapshot_value["tangle_outbound_queue_full_closes_total"], 1); 3453 assert_eq!(snapshot_value["tangle_outbox_pending_events"], 2); 3454 assert_eq!(snapshot_value["tangle_outbox_replayed_events_total"], 1); 3455 assert_eq!(snapshot_value["tangle_disk_used_bytes"], 5); 3456 assert_eq!( 3457 snapshot_value["tangle_event_admission_latency_total_micros"], 3458 13 3459 ); 3460 assert_eq!(snapshot_value["tangle_event_admission_latency_count"], 1); 3461 assert_eq!(snapshot_value["tangle_query_latency_total_micros"], 17); 3462 assert_eq!(snapshot_value["tangle_query_latency_count"], 1); 3463 assert_eq!(snapshot_value["tangle_query_candidates_scanned_total"], 5); 3464 assert_eq!(snapshot_value["tangle_query_returned_events_total"], 3); 3465 assert_eq!(snapshot_value["tangle_query_redacted_events_total"], 2); 3466 assert_eq!(snapshot_value["tangle_count_refusals_total"], 1); 3467 assert_eq!(snapshot_value["tangle_broad_query_rejections_total"], 1); 3468 3469 let report = runtime.shutdown().expect("shutdown"); 3470 3471 assert_eq!(report.closed_subscriptions(), 0); 3472 assert!(runtime.shutdown_signal().requested()); 3473 assert!(*shutdown.borrow()); 3474 3475 let _ = std::fs::remove_dir_all(root); 3476 } 3477 3478 #[test] 3479 fn runtime_metrics_snapshot_serializes_tangle_contract_keys() { 3480 let metrics = super::TangleRuntimeMetrics::new(); 3481 metrics.record_session_opened(); 3482 metrics.record_auth_success(); 3483 metrics.record_event_admission(); 3484 metrics.record_event_bus_publish(1); 3485 metrics.record_disk_used_bytes(42); 3486 let snapshot = metrics.snapshot_with_readiness(true); 3487 let value = serde_json::to_value(snapshot).expect("snapshot"); 3488 3489 assert_eq!(value["tangle_readiness_ready"], true); 3490 assert_eq!(value["tangle_ws_connections_current"], 1); 3491 assert_eq!(value["tangle_subscriptions_current"], 0); 3492 assert_eq!(value["tangle_auth_success_total"], 1); 3493 assert_eq!(value["tangle_event_admitted_total"], 1); 3494 assert_eq!(value["tangle_event_bus_published_offsets_total"], 1); 3495 assert_eq!(value["tangle_disk_used_bytes"], 42); 3496 assert_eq!(value["tangle_outbound_queue_full_closes_total"], 0); 3497 assert_eq!(value["tangle_query_candidates_scanned_total"], 0); 3498 assert_eq!(value["tangle_query_returned_events_total"], 0); 3499 assert_eq!(value["tangle_query_redacted_events_total"], 0); 3500 assert_eq!(value["tangle_count_refusals_total"], 0); 3501 assert_eq!(value["tangle_broad_query_rejections_total"], 0); 3502 assert!(value.get("active_sessions").is_none()); 3503 assert!(value.get("stored_event_offsets").is_none()); 3504 } 3505 3506 #[test] 3507 fn runtime_limits_and_event_bus_reject_zero_capacity() { 3508 assert!(TangleRuntimeLimits::new(0, runtime_relay_limits(1), 1, 1).is_err()); 3509 assert!(TangleRuntimeLimits::new(1, runtime_relay_limits(1), 0, 1).is_err()); 3510 assert!(TangleRuntimeLimits::new(1, runtime_relay_limits(1), 1, 0).is_err()); 3511 assert!(TangleEventBus::new(0).is_err()); 3512 } 3513 3514 #[tokio::test] 3515 async fn runtime_publishes_stored_event_offsets_for_live_fanout() { 3516 let root = temp_root("runtime-offset-fanout"); 3517 let _ = std::fs::remove_dir_all(&root); 3518 let handle = 3519 RelayRuntimeHandle::new(RelayRuntime::open(runtime_config(&root, 8)).expect("runtime")); 3520 let mut offsets = handle.subscribe_events().await; 3521 let mut auth = handle.auth_state().await.expect("auth"); 3522 let mut subscriptions = LiveSubscriptionSet::new(8, 64).expect("subscriptions"); 3523 let subscription_id = SubscriptionId::new("live-offset").expect("subscription"); 3524 subscriptions 3525 .subscribe( 3526 subscription_id.clone(), 3527 vec![pocket_filter(json!({"kinds":[1]}))], 3528 ) 3529 .expect("subscribe"); 3530 let event = tangle_v2_event(FixtureKey::Member, 1_714_124_433, 1, Vec::new(), "live") 3531 .expect("event"); 3532 3533 assert_eq!( 3534 handle 3535 .handle_protocol_client_message_for_test( 3536 ClientMessage::Event(event.clone()), 3537 &mut auth, 3538 UnixTimestamp::new(1_714_124_433) 3539 ) 3540 .await 3541 .expect("event"), 3542 vec![RelayMessage::Ok { 3543 event_id: event.id().clone(), 3544 accepted: true, 3545 message: String::new() 3546 }] 3547 ); 3548 let offset = offsets.try_recv().expect("offset"); 3549 assert!(matches!( 3550 handle 3551 .fanout_event_offset_with_projection_context( 3552 offset, 3553 &mut subscriptions, 3554 &auth, 3555 &RelayProjectionContext::default(), 3556 ) 3557 .await 3558 .expect("fanout") 3559 .as_slice(), 3560 [RuntimeRelayMessage::Event { 3561 subscription_id: delivered, 3562 event: found 3563 }] if delivered == &subscription_id && found.id().as_hex_string() == event.id().as_str() 3564 )); 3565 3566 assert_eq!( 3567 handle 3568 .handle_protocol_client_message_for_test( 3569 ClientMessage::Event(event.clone()), 3570 &mut auth, 3571 UnixTimestamp::new(1_714_124_434) 3572 ) 3573 .await 3574 .expect("duplicate"), 3575 vec![RelayMessage::Ok { 3576 event_id: event.id().clone(), 3577 accepted: true, 3578 message: "duplicate: already have this event".to_owned() 3579 }] 3580 ); 3581 assert_eq!( 3582 offsets.try_recv().expect_err("no duplicate offset"), 3583 TangleEventReceiveError::Empty 3584 ); 3585 let snapshot = handle.metrics().snapshot(); 3586 assert_eq!(snapshot.client_messages(), 2); 3587 assert_eq!(snapshot.event_messages(), 2); 3588 assert_eq!(snapshot.stored_event_offsets(), 1); 3589 assert_eq!(handle.metrics().event_admissions(), 2); 3590 assert_eq!(handle.metrics().event_bus_receivers_current(), 1); 3591 assert_eq!(handle.metrics().event_bus_published_offsets(), 1); 3592 assert_eq!(handle.metrics().event_admission_latency_count(), 2); 3593 3594 let _ = std::fs::remove_dir_all(root); 3595 } 3596 3597 #[tokio::test] 3598 async fn runtime_runs_stored_event_hook_before_offset_publish() { 3599 let root = temp_root("runtime-stored-hook-before-offset-publish"); 3600 let _ = std::fs::remove_dir_all(&root); 3601 let (started_sender, started_receiver) = mpsc::sync_channel(1); 3602 let (release_sender, release_receiver) = mpsc::sync_channel(1); 3603 let hooks = Arc::new(BlockingStoredHooks { 3604 started_sender: Mutex::new(Some(started_sender)), 3605 release_receiver: Mutex::new(release_receiver), 3606 }); 3607 let runtime = 3608 RelayRuntime::open_with_hooks(runtime_config(&root, 8), hooks).expect("runtime"); 3609 let handle = RelayRuntimeHandle::new(runtime); 3610 let mut offsets = handle.subscribe_events().await; 3611 let event = tangle_v2_event( 3612 FixtureKey::Member, 3613 1_714_124_433, 3614 1, 3615 Vec::new(), 3616 "stored hook before publish", 3617 ) 3618 .expect("event"); 3619 let task_handle = handle.clone(); 3620 let event_for_task = event.clone(); 3621 let task = std::thread::spawn(move || { 3622 let runtime = tokio::runtime::Builder::new_current_thread() 3623 .enable_all() 3624 .build() 3625 .expect("test runtime"); 3626 runtime.block_on(async move { 3627 let mut auth = task_handle.auth_state().await.expect("auth"); 3628 runtime_event_reply(&task_handle, event_for_task, &mut auth, 1_714_124_433).await 3629 }) 3630 }); 3631 3632 started_receiver 3633 .recv_timeout(Duration::from_secs(3)) 3634 .expect("stored hook started"); 3635 assert!(offsets.try_recv().is_err()); 3636 release_sender.send(()).expect("release hook"); 3637 assert_accepted_reply(task.join().expect("event task"), &event); 3638 offsets.try_recv().expect("offset after hook release"); 3639 handle.shutdown().await.expect("shutdown"); 3640 3641 let _ = std::fs::remove_dir_all(root); 3642 } 3643 3644 #[tokio::test] 3645 async fn runtime_hooks_reject_events_and_observe_stored_offsets() { 3646 let root = temp_root("runtime-hooks"); 3647 let _ = std::fs::remove_dir_all(&root); 3648 let hooks = Arc::new(RecordingHooks::default()); 3649 let handle = RelayRuntimeHandle::new( 3650 RelayRuntime::open_with_hooks(runtime_config(&root, 8), hooks.clone()) 3651 .expect("runtime"), 3652 ); 3653 let mut offsets = handle.subscribe_events().await; 3654 let mut auth = handle.auth_state().await.expect("auth"); 3655 let rejected = tangle_v2_event( 3656 FixtureKey::Member, 3657 1_714_124_433, 3658 1, 3659 vec![Tag::from_parts("policy", &["reject"]).expect("policy")], 3660 "rejected", 3661 ) 3662 .expect("rejected event"); 3663 3664 assert_eq!( 3665 handle 3666 .handle_protocol_client_message_for_test( 3667 ClientMessage::Event(rejected.clone()), 3668 &mut auth, 3669 UnixTimestamp::new(1_714_124_433) 3670 ) 3671 .await 3672 .expect("rejected"), 3673 vec![RelayMessage::Ok { 3674 event_id: rejected.id().clone(), 3675 accepted: false, 3676 message: "restricted: hook rejected event".to_owned() 3677 }] 3678 ); 3679 assert_eq!( 3680 offsets.try_recv().expect_err("no rejected offset"), 3681 TangleEventReceiveError::Empty 3682 ); 3683 3684 let accepted = tangle_v2_event( 3685 FixtureKey::Member, 3686 1_714_124_434, 3687 1, 3688 vec![Tag::from_parts("policy", &["accept"]).expect("policy")], 3689 "accepted", 3690 ) 3691 .expect("accepted event"); 3692 3693 assert_eq!( 3694 handle 3695 .handle_protocol_client_message_for_test( 3696 ClientMessage::Event(accepted.clone()), 3697 &mut auth, 3698 UnixTimestamp::new(1_714_124_434) 3699 ) 3700 .await 3701 .expect("accepted"), 3702 vec![RelayMessage::Ok { 3703 event_id: accepted.id().clone(), 3704 accepted: true, 3705 message: String::new() 3706 }] 3707 ); 3708 assert!(offsets.try_recv().is_ok()); 3709 let admissions = hooks.admissions.lock().expect("admissions"); 3710 assert_eq!(admissions.len(), 2); 3711 assert_eq!(admissions[0].event().event_id(), rejected.id().as_str()); 3712 assert_eq!(admissions[0].event().created_at(), 1_714_124_433); 3713 assert_eq!(admissions[1].event().event_id(), accepted.id().as_str()); 3714 assert_eq!(admissions[1].event().created_at(), 1_714_124_434); 3715 drop(admissions); 3716 let stored = hooks.stored.lock().expect("stored"); 3717 assert_eq!(stored.len(), 1); 3718 assert_eq!(stored[0].event().event_id(), accepted.id().as_str()); 3719 assert_eq!(stored[0].event().created_at(), 1_714_124_434); 3720 assert_eq!(stored[0].store_offsets().len(), 1); 3721 assert_eq!(handle.metrics().disk_used_bytes(), 144); 3722 3723 let _ = std::fs::remove_dir_all(root); 3724 } 3725 3726 #[tokio::test] 3727 async fn runtime_projection_default_keeps_query_and_live_output() { 3728 let root = temp_root("runtime-projection-default"); 3729 let _ = std::fs::remove_dir_all(&root); 3730 let handle = 3731 RelayRuntimeHandle::new(RelayRuntime::open(runtime_config(&root, 8)).expect("runtime")); 3732 let mut offsets = handle.subscribe_events().await; 3733 let mut auth = handle.auth_state().await.expect("auth"); 3734 let event = tangle_v2_event( 3735 FixtureKey::Member, 3736 1_714_124_433, 3737 1, 3738 Vec::new(), 3739 "default projection", 3740 ) 3741 .expect("event"); 3742 3743 assert_eq!( 3744 handle 3745 .handle_protocol_client_message_for_test( 3746 ClientMessage::Event(event.clone()), 3747 &mut auth, 3748 UnixTimestamp::new(1_714_124_433) 3749 ) 3750 .await 3751 .expect("event"), 3752 vec![RelayMessage::Ok { 3753 event_id: event.id().clone(), 3754 accepted: true, 3755 message: String::new() 3756 }] 3757 ); 3758 let offset = offsets.try_recv().expect("offset"); 3759 let query_sub = SubscriptionId::new("projection-default-query").expect("subscription"); 3760 let report = handle 3761 .query_req_with_auth_report_with_projection_context( 3762 query_sub.clone(), 3763 vec![pocket_filter(json!({"ids": [event.id().as_str()]}))], 3764 false, 3765 &auth, 3766 &RelayProjectionContext::default(), 3767 ) 3768 .await 3769 .expect("query"); 3770 assert!(matches!( 3771 report.into_messages().as_slice(), 3772 [ 3773 RuntimeRelayMessage::Event { 3774 subscription_id, 3775 event: found 3776 }, 3777 RuntimeRelayMessage::Protocol(RelayMessage::Eose(eose)) 3778 ] if subscription_id == &query_sub 3779 && found.id().as_hex_string() == event.id().as_str() 3780 && eose == &query_sub 3781 )); 3782 3783 let live_sub = SubscriptionId::new("projection-default-live").expect("subscription"); 3784 let mut subscriptions = LiveSubscriptionSet::new(8, 64).expect("subscriptions"); 3785 subscriptions 3786 .subscribe(live_sub.clone(), vec![pocket_filter(json!({"kinds": [1]}))]) 3787 .expect("subscribe"); 3788 assert!(matches!( 3789 handle 3790 .fanout_event_offset_with_projection_context( 3791 offset, 3792 &mut subscriptions, 3793 &auth, 3794 &RelayProjectionContext::default() 3795 ) 3796 .await 3797 .expect("fanout") 3798 .as_slice(), 3799 [RuntimeRelayMessage::Event { 3800 subscription_id, 3801 event: found 3802 }] if subscription_id == &live_sub && found.id().as_hex_string() == event.id().as_str() 3803 )); 3804 3805 let _ = std::fs::remove_dir_all(root); 3806 } 3807 3808 #[tokio::test] 3809 async fn runtime_projection_can_suppress_historical_query_events() { 3810 let root = temp_root("runtime-projection-query-suppress"); 3811 let _ = std::fs::remove_dir_all(&root); 3812 let hooks = Arc::new(ProjectingHooks::new( 3813 "quiet", 3814 ProjectionHookScope::Historical, 3815 None, 3816 RelayEventProjectionDecision::Suppress, 3817 )); 3818 let handle = RelayRuntimeHandle::new( 3819 RelayRuntime::open_with_hooks(runtime_config(&root, 8), hooks.clone()) 3820 .expect("runtime"), 3821 ); 3822 let mut auth = handle.auth_state().await.expect("auth"); 3823 let event = tangle_v2_event( 3824 FixtureKey::Member, 3825 1_714_124_433, 3826 1, 3827 Vec::new(), 3828 "suppress query", 3829 ) 3830 .expect("event"); 3831 assert_accepted_reply( 3832 runtime_event_reply(&handle, event.clone(), &mut auth, 1_714_124_433).await, 3833 &event, 3834 ); 3835 let subscription_id = SubscriptionId::new("query-suppressed").expect("subscription"); 3836 3837 let report = handle 3838 .query_req_with_auth_report_with_projection_context( 3839 subscription_id.clone(), 3840 vec![pocket_filter(json!({"ids": [event.id().as_str()]}))], 3841 false, 3842 &auth, 3843 &RelayProjectionContext::named("quiet").expect("projection"), 3844 ) 3845 .await 3846 .expect("query"); 3847 assert_eq!( 3848 report.into_messages(), 3849 vec![RuntimeRelayMessage::from(RelayMessage::Eose( 3850 subscription_id 3851 ))] 3852 ); 3853 let contexts = hooks.contexts(); 3854 assert_eq!(contexts.len(), 1); 3855 assert_eq!(contexts[0].projection().identifier(), Some("quiet")); 3856 assert_eq!( 3857 contexts[0].source(), 3858 RelayEventProjectionSource::HistoricalQuery 3859 ); 3860 assert_eq!(contexts[0].matched_filter().filter_index(), 0); 3861 assert_eq!( 3862 contexts[0].matched_filter().requested_kinds(), 3863 &RelayRequestedKinds::Absent 3864 ); 3865 assert_eq!(contexts[0].event().event_id(), event.id().as_str()); 3866 let query_contexts = hooks.query_contexts(); 3867 assert_eq!(query_contexts.len(), 1); 3868 assert_eq!( 3869 query_contexts[0].filters()[0].requested_kinds(), 3870 &RelayRequestedKinds::Absent 3871 ); 3872 3873 let _ = std::fs::remove_dir_all(root); 3874 } 3875 3876 #[tokio::test] 3877 async fn runtime_projection_context_reports_all_matching_query_filters() { 3878 let root = temp_root("runtime-projection-query-all-matched-filters"); 3879 let _ = std::fs::remove_dir_all(&root); 3880 let hooks = Arc::new(ProjectingHooks::new( 3881 "multi-filter", 3882 ProjectionHookScope::Historical, 3883 None, 3884 RelayEventProjectionDecision::Emit, 3885 )); 3886 let handle = RelayRuntimeHandle::new( 3887 RelayRuntime::open_with_hooks(runtime_config(&root, 8), hooks.clone()) 3888 .expect("runtime"), 3889 ); 3890 let mut auth = handle.auth_state().await.expect("auth"); 3891 let event = tangle_v2_event( 3892 FixtureKey::Member, 3893 1_714_124_433, 3894 1, 3895 Vec::new(), 3896 "multi filter query", 3897 ) 3898 .expect("event"); 3899 assert_accepted_reply( 3900 runtime_event_reply(&handle, event.clone(), &mut auth, 1_714_124_433).await, 3901 &event, 3902 ); 3903 let subscription_id = SubscriptionId::new("query-multi-filter").expect("subscription"); 3904 3905 let report = handle 3906 .query_req_with_auth_report_with_projection_context( 3907 subscription_id.clone(), 3908 vec![ 3909 pocket_filter(json!({"ids": [event.id().as_str()]})), 3910 pocket_filter(json!({"kinds": [1]})), 3911 ], 3912 false, 3913 &auth, 3914 &RelayProjectionContext::named("multi-filter").expect("projection"), 3915 ) 3916 .await 3917 .expect("query"); 3918 assert!(matches!( 3919 report.into_messages().as_slice(), 3920 [ 3921 RuntimeRelayMessage::Event { 3922 subscription_id: delivered, 3923 event: delivered_event 3924 }, 3925 RuntimeRelayMessage::Protocol(RelayMessage::Eose(eose)) 3926 ] if delivered == &subscription_id 3927 && delivered_event.id().as_hex_string() == event.id().as_str() 3928 && eose == &subscription_id 3929 )); 3930 let contexts = hooks.contexts(); 3931 assert_eq!(contexts.len(), 1); 3932 assert_eq!(contexts[0].matched_filter().filter_index(), 0); 3933 let matched_filters = contexts[0].matched_filters(); 3934 assert_eq!(matched_filters.len(), 2); 3935 assert_eq!(matched_filters[0].filter_index(), 0); 3936 assert_eq!( 3937 matched_filters[0].requested_kinds(), 3938 &RelayRequestedKinds::Absent 3939 ); 3940 assert_eq!(matched_filters[1].filter_index(), 1); 3941 assert_eq!( 3942 matched_filters[1].requested_kinds(), 3943 &RelayRequestedKinds::Explicit(BTreeSet::from([1])) 3944 ); 3945 3946 let _ = std::fs::remove_dir_all(root); 3947 } 3948 3949 #[tokio::test] 3950 async fn runtime_projection_post_limit_reports_all_matching_query_filters() { 3951 let root = temp_root("runtime-projection-post-limit-all-matched-filters"); 3952 let _ = std::fs::remove_dir_all(&root); 3953 let hooks = Arc::new(ProjectingHooks::new( 3954 "post-limit-multi-filter", 3955 ProjectionHookScope::Historical, 3956 None, 3957 RelayEventProjectionDecision::Emit, 3958 )); 3959 hooks.set_query_plan(RelayProjectionQueryPlan::limit_after_projection( 3960 NonZeroU32::new(10).expect("candidate limit"), 3961 )); 3962 let handle = RelayRuntimeHandle::new( 3963 RelayRuntime::open_with_hooks(runtime_config(&root, 8), hooks.clone()) 3964 .expect("runtime"), 3965 ); 3966 let mut auth = handle.auth_state().await.expect("auth"); 3967 let event = tangle_v2_event( 3968 FixtureKey::Member, 3969 1_714_124_433, 3970 30078, 3971 Vec::new(), 3972 "post limit multi filter query", 3973 ) 3974 .expect("event"); 3975 assert_accepted_reply( 3976 runtime_event_reply(&handle, event.clone(), &mut auth, 1_714_124_433).await, 3977 &event, 3978 ); 3979 let subscription_id = 3980 SubscriptionId::new("post-limit-query-multi-filter").expect("subscription"); 3981 3982 let report = handle 3983 .query_req_with_auth_report_with_projection_context( 3984 subscription_id.clone(), 3985 vec![ 3986 pocket_filter(json!({"kinds": [30078]})), 3987 pocket_filter(json!({"ids": [event.id().as_str()]})), 3988 ], 3989 false, 3990 &auth, 3991 &RelayProjectionContext::named("post-limit-multi-filter").expect("projection"), 3992 ) 3993 .await 3994 .expect("query"); 3995 assert!(matches!( 3996 report.into_messages().as_slice(), 3997 [ 3998 RuntimeRelayMessage::Event { 3999 subscription_id: delivered, 4000 event: delivered_event 4001 }, 4002 RuntimeRelayMessage::Protocol(RelayMessage::Eose(eose)) 4003 ] if delivered == &subscription_id 4004 && delivered_event.id().as_hex_string() == event.id().as_str() 4005 && eose == &subscription_id 4006 )); 4007 let contexts = hooks.contexts(); 4008 assert_eq!(contexts.len(), 2); 4009 for context in contexts { 4010 let matched_filters = context.matched_filters(); 4011 assert_eq!(matched_filters.len(), 2); 4012 assert_eq!(matched_filters[0].filter_index(), 0); 4013 assert_eq!( 4014 matched_filters[0].requested_kinds(), 4015 &RelayRequestedKinds::Explicit(BTreeSet::from([30078])) 4016 ); 4017 assert_eq!(matched_filters[1].filter_index(), 1); 4018 assert_eq!( 4019 matched_filters[1].requested_kinds(), 4020 &RelayRequestedKinds::Absent 4021 ); 4022 } 4023 4024 let _ = std::fs::remove_dir_all(root); 4025 } 4026 4027 #[tokio::test] 4028 async fn runtime_projection_can_suppress_live_fanout_events() { 4029 let root = temp_root("runtime-projection-live-suppress"); 4030 let _ = std::fs::remove_dir_all(&root); 4031 let hooks = Arc::new(ProjectingHooks::new( 4032 "quiet-live", 4033 ProjectionHookScope::Live, 4034 None, 4035 RelayEventProjectionDecision::Suppress, 4036 )); 4037 let handle = RelayRuntimeHandle::new( 4038 RelayRuntime::open_with_hooks(runtime_config(&root, 8), hooks.clone()) 4039 .expect("runtime"), 4040 ); 4041 let mut offsets = handle.subscribe_events().await; 4042 let mut auth = handle.auth_state().await.expect("auth"); 4043 let event = tangle_v2_event( 4044 FixtureKey::Member, 4045 1_714_124_433, 4046 1, 4047 Vec::new(), 4048 "suppress live", 4049 ) 4050 .expect("event"); 4051 assert_accepted_reply( 4052 runtime_event_reply(&handle, event.clone(), &mut auth, 1_714_124_433).await, 4053 &event, 4054 ); 4055 let offset = offsets.try_recv().expect("offset"); 4056 let mut subscriptions = LiveSubscriptionSet::new(8, 64).expect("subscriptions"); 4057 subscriptions 4058 .subscribe( 4059 SubscriptionId::new("live-suppressed").expect("subscription"), 4060 vec![pocket_filter(json!({"kinds": [1]}))], 4061 ) 4062 .expect("subscribe"); 4063 4064 assert!( 4065 handle 4066 .fanout_event_offset_with_projection_context( 4067 offset, 4068 &mut subscriptions, 4069 &auth, 4070 &RelayProjectionContext::named("quiet-live").expect("projection") 4071 ) 4072 .await 4073 .expect("fanout") 4074 .is_empty() 4075 ); 4076 let contexts = hooks.contexts(); 4077 assert_eq!(contexts.len(), 1); 4078 assert_eq!(contexts[0].projection().identifier(), Some("quiet-live")); 4079 assert_eq!( 4080 contexts[0].source(), 4081 RelayEventProjectionSource::LiveFanout { 4082 store_offset: offset.as_u64() 4083 } 4084 ); 4085 assert_eq!(contexts[0].matched_filter().filter_index(), 0); 4086 assert_eq!( 4087 contexts[0].matched_filter().requested_kinds(), 4088 &RelayRequestedKinds::Explicit(BTreeSet::from([1])) 4089 ); 4090 assert_eq!(contexts[0].event().event_id(), event.id().as_str()); 4091 4092 let _ = std::fs::remove_dir_all(root); 4093 } 4094 4095 #[tokio::test] 4096 async fn runtime_projection_contexts_include_authenticated_pubkeys() { 4097 let root = temp_root("runtime-projection-authenticated-pubkeys"); 4098 let _ = std::fs::remove_dir_all(&root); 4099 let hooks = Arc::new(ProjectingHooks::new( 4100 "auth-context", 4101 ProjectionHookScope::Historical, 4102 None, 4103 RelayEventProjectionDecision::Emit, 4104 )); 4105 let handle = RelayRuntimeHandle::new( 4106 RelayRuntime::open_with_hooks(runtime_config(&root, 8), hooks.clone()) 4107 .expect("runtime"), 4108 ); 4109 let mut offsets = handle.subscribe_events().await; 4110 let mut auth = 4111 authenticated_runtime_state(&handle, FixtureKey::Owner, "challenge-auth-context", 100) 4112 .await; 4113 let expected_pubkeys = auth 4114 .authenticated_pubkeys() 4115 .iter() 4116 .cloned() 4117 .collect::<Vec<_>>(); 4118 let event = tangle_v2_event( 4119 FixtureKey::Member, 4120 1_714_124_433, 4121 1, 4122 Vec::new(), 4123 "authenticated projection context", 4124 ) 4125 .expect("event"); 4126 assert_accepted_reply( 4127 runtime_event_reply(&handle, event.clone(), &mut auth, 1_714_124_433).await, 4128 &event, 4129 ); 4130 let offset = offsets.try_recv().expect("offset"); 4131 let query_sub = SubscriptionId::new("query-auth-context").expect("subscription"); 4132 let projection = RelayProjectionContext::named("auth-context").expect("projection"); 4133 4134 let report = handle 4135 .query_req_with_auth_report_with_projection_context( 4136 query_sub.clone(), 4137 vec![pocket_filter(json!({"ids": [event.id().as_str()]}))], 4138 false, 4139 &auth, 4140 &projection, 4141 ) 4142 .await 4143 .expect("query"); 4144 assert!(matches!( 4145 report.into_messages().as_slice(), 4146 [ 4147 RuntimeRelayMessage::Event { 4148 subscription_id, 4149 event: found 4150 }, 4151 RuntimeRelayMessage::Protocol(RelayMessage::Eose(eose)) 4152 ] if subscription_id == &query_sub 4153 && found.id().as_hex_string() == event.id().as_str() 4154 && eose == &query_sub 4155 )); 4156 4157 let mut subscriptions = LiveSubscriptionSet::new(8, 64).expect("subscriptions"); 4158 subscriptions 4159 .subscribe( 4160 SubscriptionId::new("live-auth-context").expect("subscription"), 4161 vec![pocket_filter(json!({"kinds": [1]}))], 4162 ) 4163 .expect("subscribe"); 4164 assert_eq!( 4165 handle 4166 .fanout_event_offset_with_projection_context( 4167 offset, 4168 &mut subscriptions, 4169 &auth, 4170 &projection, 4171 ) 4172 .await 4173 .expect("fanout") 4174 .len(), 4175 1 4176 ); 4177 4178 let query_contexts = hooks.query_contexts(); 4179 assert_eq!(query_contexts.len(), 1); 4180 assert_eq!( 4181 query_contexts[0].authenticated_pubkeys(), 4182 expected_pubkeys.as_slice() 4183 ); 4184 let live_contexts = hooks.live_contexts(); 4185 assert_eq!(live_contexts.len(), 1); 4186 assert_eq!( 4187 live_contexts[0].authenticated_pubkeys(), 4188 expected_pubkeys.as_slice() 4189 ); 4190 let contexts = hooks.contexts(); 4191 assert_eq!(contexts.len(), 2); 4192 assert_eq!( 4193 contexts[0].source(), 4194 RelayEventProjectionSource::HistoricalQuery 4195 ); 4196 assert_eq!( 4197 contexts[0].authenticated_pubkeys(), 4198 expected_pubkeys.as_slice() 4199 ); 4200 assert_eq!( 4201 contexts[1].source(), 4202 RelayEventProjectionSource::LiveFanout { 4203 store_offset: offset.as_u64() 4204 } 4205 ); 4206 assert_eq!( 4207 contexts[1].authenticated_pubkeys(), 4208 expected_pubkeys.as_slice() 4209 ); 4210 4211 let _ = std::fs::remove_dir_all(root); 4212 } 4213 4214 #[tokio::test] 4215 async fn runtime_projection_replaces_with_existing_stored_events_only() { 4216 let root = temp_root("runtime-projection-replace"); 4217 let _ = std::fs::remove_dir_all(&root); 4218 let hooks = Arc::new(ProjectingHooks::new( 4219 "replace", 4220 ProjectionHookScope::Historical, 4221 Some("source"), 4222 RelayEventProjectionDecision::Emit, 4223 )); 4224 let handle = RelayRuntimeHandle::new( 4225 RelayRuntime::open_with_hooks(runtime_config(&root, 8), hooks.clone()) 4226 .expect("runtime"), 4227 ); 4228 let mut offsets = handle.subscribe_events().await; 4229 let mut auth = handle.auth_state().await.expect("auth"); 4230 let source = tangle_v2_event(FixtureKey::Member, 1_714_124_433, 1, Vec::new(), "source") 4231 .expect("source"); 4232 let replacement = tangle_v2_event( 4233 FixtureKey::Admin, 4234 1_714_124_434, 4235 1, 4236 Vec::new(), 4237 "replacement", 4238 ) 4239 .expect("replacement"); 4240 assert_accepted_reply( 4241 runtime_event_reply(&handle, source.clone(), &mut auth, 1_714_124_433).await, 4242 &source, 4243 ); 4244 let _source_offset = offsets.try_recv().expect("source offset"); 4245 assert_accepted_reply( 4246 runtime_event_reply(&handle, replacement.clone(), &mut auth, 1_714_124_434).await, 4247 &replacement, 4248 ); 4249 let replacement_offset = offsets.try_recv().expect("replacement offset"); 4250 hooks.set_decision(RelayEventProjectionDecision::replace_with_stored_offset( 4251 replacement_offset.as_u64(), 4252 )); 4253 let subscription_id = SubscriptionId::new("replace-existing").expect("subscription"); 4254 4255 let report = handle 4256 .query_req_with_auth_report_with_projection_context( 4257 subscription_id.clone(), 4258 vec![pocket_filter(json!({"kinds": [1]}))], 4259 false, 4260 &auth, 4261 &RelayProjectionContext::named("replace").expect("projection"), 4262 ) 4263 .await 4264 .expect("query"); 4265 assert!(matches!( 4266 report.into_messages().as_slice(), 4267 [ 4268 RuntimeRelayMessage::Event { 4269 subscription_id: delivered, 4270 event 4271 }, 4272 RuntimeRelayMessage::Protocol(RelayMessage::Eose(eose)) 4273 ] if delivered == &subscription_id 4274 && event.id().as_hex_string() == replacement.id().as_str() 4275 && eose == &subscription_id 4276 )); 4277 4278 hooks.set_decision(RelayEventProjectionDecision::replace_with_stored_offset( 4279 u64::MAX, 4280 )); 4281 let missing_subscription = SubscriptionId::new("replace-missing").expect("subscription"); 4282 let report = handle 4283 .query_req_with_auth_report_with_projection_context( 4284 missing_subscription.clone(), 4285 vec![pocket_filter(json!({"ids": [source.id().as_str()]}))], 4286 false, 4287 &auth, 4288 &RelayProjectionContext::named("replace").expect("projection"), 4289 ) 4290 .await 4291 .expect("missing query"); 4292 assert_eq!( 4293 report.into_messages(), 4294 vec![RuntimeRelayMessage::from(RelayMessage::Eose( 4295 missing_subscription 4296 ))] 4297 ); 4298 4299 let _ = std::fs::remove_dir_all(root); 4300 } 4301 4302 #[tokio::test] 4303 async fn runtime_projection_can_apply_query_limit_after_projection() { 4304 let root = temp_root("runtime-projection-post-limit"); 4305 let _ = std::fs::remove_dir_all(&root); 4306 let hooks = Arc::new(ProjectingHooks::new( 4307 "post-limit", 4308 ProjectionHookScope::Historical, 4309 Some("drop"), 4310 RelayEventProjectionDecision::Suppress, 4311 )); 4312 hooks.set_query_plan(RelayProjectionQueryPlan::limit_after_projection( 4313 NonZeroU32::new(2).expect("candidate limit"), 4314 )); 4315 let handle = RelayRuntimeHandle::new( 4316 RelayRuntime::open_with_hooks(runtime_config(&root, 8), hooks.clone()) 4317 .expect("runtime"), 4318 ); 4319 let mut auth = handle.auth_state().await.expect("auth"); 4320 let dropped = tangle_v2_event(FixtureKey::Member, 1_714_124_435, 1, Vec::new(), "drop") 4321 .expect("drop"); 4322 let kept = 4323 tangle_v2_event(FixtureKey::Admin, 1_714_124_434, 1, Vec::new(), "keep").expect("keep"); 4324 assert_accepted_reply( 4325 runtime_event_reply(&handle, kept.clone(), &mut auth, 1_714_124_434).await, 4326 &kept, 4327 ); 4328 assert_accepted_reply( 4329 runtime_event_reply(&handle, dropped.clone(), &mut auth, 1_714_124_435).await, 4330 &dropped, 4331 ); 4332 let subscription_id = SubscriptionId::new("post-limit").expect("subscription"); 4333 4334 let report = handle 4335 .query_req_with_auth_report_with_projection_context( 4336 subscription_id.clone(), 4337 vec![pocket_filter(json!({"kinds": [1], "limit": 1}))], 4338 false, 4339 &auth, 4340 &RelayProjectionContext::named("post-limit").expect("projection"), 4341 ) 4342 .await 4343 .expect("query"); 4344 assert!(matches!( 4345 report.into_messages().as_slice(), 4346 [ 4347 RuntimeRelayMessage::Event { 4348 subscription_id: delivered, 4349 event 4350 }, 4351 RuntimeRelayMessage::Protocol(RelayMessage::Eose(eose)) 4352 ] if delivered == &subscription_id 4353 && event.id().as_hex_string() == kept.id().as_str() 4354 && eose == &subscription_id 4355 )); 4356 let query_contexts = hooks.query_contexts(); 4357 assert_eq!(query_contexts.len(), 1); 4358 assert_eq!( 4359 query_contexts[0].filters()[0].requested_kinds(), 4360 &RelayRequestedKinds::Explicit(BTreeSet::from([1])) 4361 ); 4362 assert_eq!(hooks.contexts().len(), 2); 4363 4364 let _ = std::fs::remove_dir_all(root); 4365 } 4366 4367 #[tokio::test] 4368 async fn runtime_projection_live_replacement_must_match_original_filter() { 4369 let root = temp_root("runtime-projection-live-replace-rematch"); 4370 let _ = std::fs::remove_dir_all(&root); 4371 let hooks = Arc::new(ProjectingHooks::new( 4372 "live-replace", 4373 ProjectionHookScope::Live, 4374 Some("source"), 4375 RelayEventProjectionDecision::Emit, 4376 )); 4377 let handle = RelayRuntimeHandle::new( 4378 RelayRuntime::open_with_hooks(runtime_config(&root, 8), hooks.clone()) 4379 .expect("runtime"), 4380 ); 4381 let mut offsets = handle.subscribe_events().await; 4382 let mut auth = handle.auth_state().await.expect("auth"); 4383 let source = tangle_v2_event(FixtureKey::Member, 1_714_124_433, 7, Vec::new(), "source") 4384 .expect("source"); 4385 let replacement = tangle_v2_event( 4386 FixtureKey::Admin, 4387 1_714_124_434, 4388 1, 4389 Vec::new(), 4390 "replacement", 4391 ) 4392 .expect("replacement"); 4393 assert_accepted_reply( 4394 runtime_event_reply(&handle, source.clone(), &mut auth, 1_714_124_433).await, 4395 &source, 4396 ); 4397 let source_offset = offsets.try_recv().expect("source offset"); 4398 assert_accepted_reply( 4399 runtime_event_reply(&handle, replacement.clone(), &mut auth, 1_714_124_434).await, 4400 &replacement, 4401 ); 4402 let replacement_offset = offsets.try_recv().expect("replacement offset"); 4403 hooks.set_decision(RelayEventProjectionDecision::replace_with_stored_offset( 4404 replacement_offset.as_u64(), 4405 )); 4406 let mut subscriptions = LiveSubscriptionSet::new(8, 64).expect("subscriptions"); 4407 let source_only = SubscriptionId::new("source-only").expect("source subscription"); 4408 let source_or_note = SubscriptionId::new("source-or-note").expect("mixed subscription"); 4409 subscriptions 4410 .subscribe(source_only, vec![pocket_filter(json!({"kinds": [7]}))]) 4411 .expect("source subscribe"); 4412 subscriptions 4413 .subscribe( 4414 source_or_note.clone(), 4415 vec![pocket_filter(json!({"kinds": [1, 7]}))], 4416 ) 4417 .expect("mixed subscribe"); 4418 4419 let messages = handle 4420 .fanout_event_offset_with_projection_context( 4421 source_offset, 4422 &mut subscriptions, 4423 &auth, 4424 &RelayProjectionContext::named("live-replace").expect("projection"), 4425 ) 4426 .await 4427 .expect("fanout"); 4428 assert!(matches!( 4429 messages.as_slice(), 4430 [RuntimeRelayMessage::Event { 4431 subscription_id, 4432 event 4433 }] if subscription_id == &source_or_note 4434 && event.id().as_hex_string() == replacement.id().as_str() 4435 )); 4436 let contexts = hooks.contexts(); 4437 assert_eq!(contexts.len(), 2); 4438 assert_eq!( 4439 contexts[0].matched_filter().requested_kinds(), 4440 &RelayRequestedKinds::Explicit(BTreeSet::from([7])) 4441 ); 4442 assert_eq!( 4443 contexts[1].matched_filter().requested_kinds(), 4444 &RelayRequestedKinds::Explicit(BTreeSet::from([1, 7])) 4445 ); 4446 4447 let _ = std::fs::remove_dir_all(root); 4448 } 4449 4450 #[tokio::test] 4451 async fn runtime_projection_live_candidates_match_candidate_filter() { 4452 let root = temp_root("runtime-projection-live-candidate"); 4453 let _ = std::fs::remove_dir_all(&root); 4454 let hooks = Arc::new(ProjectingHooks::new( 4455 "live-candidate", 4456 ProjectionHookScope::Live, 4457 None, 4458 RelayEventProjectionDecision::Emit, 4459 )); 4460 let handle = RelayRuntimeHandle::new( 4461 RelayRuntime::open_with_hooks(runtime_config(&root, 8), hooks.clone()) 4462 .expect("runtime"), 4463 ); 4464 let mut offsets = handle.subscribe_events().await; 4465 let mut auth = handle.auth_state().await.expect("auth"); 4466 let source = tangle_v2_event( 4467 FixtureKey::Member, 4468 1_714_124_433, 4469 30078, 4470 Vec::new(), 4471 "source", 4472 ) 4473 .expect("source"); 4474 let candidate = 4475 tangle_v2_event(FixtureKey::Admin, 1_714_124_434, 1, Vec::new(), "candidate") 4476 .expect("candidate"); 4477 assert_accepted_reply( 4478 runtime_event_reply(&handle, source.clone(), &mut auth, 1_714_124_433).await, 4479 &source, 4480 ); 4481 let source_offset = offsets.try_recv().expect("source offset"); 4482 assert_accepted_reply( 4483 runtime_event_reply(&handle, candidate.clone(), &mut auth, 1_714_124_434).await, 4484 &candidate, 4485 ); 4486 let candidate_offset = offsets.try_recv().expect("candidate offset"); 4487 hooks.set_live_candidates(vec![RelayLiveProjectionCandidate::stored_offset( 4488 candidate_offset.as_u64(), 4489 )]); 4490 let mut subscriptions = LiveSubscriptionSet::new(8, 64).expect("subscriptions"); 4491 let subscription_id = SubscriptionId::new("candidate-note").expect("subscription"); 4492 subscriptions 4493 .subscribe( 4494 subscription_id.clone(), 4495 vec![pocket_filter(json!({"kinds": [1]}))], 4496 ) 4497 .expect("subscribe"); 4498 4499 let messages = handle 4500 .fanout_event_offset_with_projection_context( 4501 source_offset, 4502 &mut subscriptions, 4503 &auth, 4504 &RelayProjectionContext::named("live-candidate").expect("projection"), 4505 ) 4506 .await 4507 .expect("fanout"); 4508 4509 assert!(matches!( 4510 messages.as_slice(), 4511 [RuntimeRelayMessage::Event { 4512 subscription_id: delivered, 4513 event 4514 }] if delivered == &subscription_id 4515 && event.id().as_hex_string() == candidate.id().as_str() 4516 )); 4517 let live_contexts = hooks.live_contexts(); 4518 assert_eq!(live_contexts.len(), 1); 4519 assert_eq!( 4520 live_contexts[0].source_store_offset(), 4521 source_offset.as_u64() 4522 ); 4523 assert_eq!(live_contexts[0].event().event_id(), source.id().as_str()); 4524 let contexts = hooks.contexts(); 4525 assert_eq!(contexts.len(), 1); 4526 assert_eq!(contexts[0].event().event_id(), candidate.id().as_str()); 4527 assert_eq!( 4528 contexts[0].matched_filter().requested_kinds(), 4529 &RelayRequestedKinds::Explicit(BTreeSet::from([1])) 4530 ); 4531 4532 let _ = std::fs::remove_dir_all(root); 4533 } 4534 4535 #[tokio::test] 4536 async fn runtime_projection_runs_after_group_read_gates() { 4537 let root = temp_root("runtime-projection-group-gate"); 4538 let _ = std::fs::remove_dir_all(&root); 4539 let hooks = Arc::new(ProjectingHooks::new( 4540 "group-gate", 4541 ProjectionHookScope::Historical, 4542 None, 4543 RelayEventProjectionDecision::Suppress, 4544 )); 4545 let handle = RelayRuntimeHandle::new( 4546 RelayRuntime::open_with_hooks(runtime_config(&root, 8), hooks.clone()) 4547 .expect("runtime"), 4548 ); 4549 let mut owner_auth = 4550 authenticated_runtime_state(&handle, FixtureKey::Owner, "group-gate-owner", 120).await; 4551 let public_auth = handle.auth_state().await.expect("public auth"); 4552 let create = 4553 tangle_v2_group_create_event(FixtureKey::Owner, "ProjectionPrivate", 121, &["private"]) 4554 .expect("create"); 4555 assert_accepted_reply( 4556 runtime_event_reply(&handle, create.clone(), &mut owner_auth, 121).await, 4557 &create, 4558 ); 4559 let private_event = tangle_v2_group_event( 4560 FixtureKey::Owner, 4561 "ProjectionPrivate", 4562 122, 4563 1, 4564 "private projection", 4565 ) 4566 .expect("private event"); 4567 assert_accepted_reply( 4568 runtime_event_reply(&handle, private_event.clone(), &mut owner_auth, 122).await, 4569 &private_event, 4570 ); 4571 let subscription_id = SubscriptionId::new("group-gate").expect("subscription"); 4572 4573 let report = handle 4574 .query_req_with_auth_report_with_projection_context( 4575 subscription_id.clone(), 4576 vec![pocket_filter(json!({ 4577 "kinds": [1], 4578 "#h": ["ProjectionPrivate"] 4579 }))], 4580 false, 4581 &public_auth, 4582 &RelayProjectionContext::named("group-gate").expect("projection"), 4583 ) 4584 .await 4585 .expect("query"); 4586 assert_eq!( 4587 report.into_messages(), 4588 vec![RuntimeRelayMessage::from(RelayMessage::Closed { 4589 subscription_id, 4590 message: "auth-required: authentication required to read group events".to_owned() 4591 })] 4592 ); 4593 assert!(hooks.contexts().is_empty()); 4594 4595 let _ = std::fs::remove_dir_all(root); 4596 } 4597 4598 #[tokio::test] 4599 async fn runtime_rate_limits_event_pubkeys_before_storage() { 4600 let root = temp_root("runtime-event-rate-limit"); 4601 let _ = std::fs::remove_dir_all(&root); 4602 let runtime = RelayRuntime::open(runtime_config(&root, 8)).expect("runtime"); 4603 let event = tangle_v2_event(FixtureKey::Member, 1_714_124_433, 1, Vec::new(), "limited") 4604 .expect("event"); 4605 let rule = runtime.config().rate_limits().event().per_pubkey(); 4606 let key = TangleRateLimitKey::pubkey( 4607 TangleRateLimitScope::Event, 4608 event.unsigned().pubkey().clone(), 4609 ); 4610 for _ in 0..rule.max_hits() { 4611 runtime 4612 .rate_limiter() 4613 .record(key.clone(), rule, UnixTimestamp::new(1_714_124_433)); 4614 } 4615 let handle = RelayRuntimeHandle::new(runtime); 4616 let mut auth = handle.auth_state().await.expect("auth"); 4617 4618 assert_eq!( 4619 handle 4620 .handle_protocol_client_message_for_test( 4621 ClientMessage::Event(event.clone()), 4622 &mut auth, 4623 UnixTimestamp::new(1_714_124_433) 4624 ) 4625 .await 4626 .expect("event"), 4627 vec![RelayMessage::Ok { 4628 event_id: event.id().clone(), 4629 accepted: false, 4630 message: "rate-limited: event pubkey rate limit exceeded until 1714124493" 4631 .to_owned() 4632 }] 4633 ); 4634 4635 let _ = std::fs::remove_dir_all(root); 4636 } 4637 4638 #[tokio::test] 4639 async fn runtime_rate_limits_event_kinds_before_storage() { 4640 let root = temp_root("runtime-event-kind-rate-limit"); 4641 let _ = std::fs::remove_dir_all(&root); 4642 let runtime = RelayRuntime::open(runtime_config(&root, 8)).expect("runtime"); 4643 let event = tangle_v2_event(FixtureKey::Admin, 1_714_124_433, 1, Vec::new(), "limited") 4644 .expect("event"); 4645 let rule = runtime.config().rate_limits().event().per_kind(); 4646 let key = TangleRateLimitKey::kind(TangleRateLimitScope::Event, event.unsigned().kind()); 4647 for _ in 0..rule.max_hits() { 4648 runtime 4649 .rate_limiter() 4650 .record(key.clone(), rule, UnixTimestamp::new(1_714_124_433)); 4651 } 4652 let handle = RelayRuntimeHandle::new(runtime); 4653 let mut auth = handle.auth_state().await.expect("auth"); 4654 4655 assert_eq!( 4656 handle 4657 .handle_protocol_client_message_for_test( 4658 ClientMessage::Event(event.clone()), 4659 &mut auth, 4660 UnixTimestamp::new(1_714_124_433) 4661 ) 4662 .await 4663 .expect("event"), 4664 vec![RelayMessage::Ok { 4665 event_id: event.id().clone(), 4666 accepted: false, 4667 message: "rate-limited: event kind rate limit exceeded until 1714124493".to_owned() 4668 }] 4669 ); 4670 4671 let _ = std::fs::remove_dir_all(root); 4672 } 4673 4674 #[tokio::test] 4675 async fn runtime_rate_limits_event_peer_ips_partition_peers_and_precede_identity_keys() { 4676 let root = temp_root("runtime-event-ip-rate-limit"); 4677 let _ = std::fs::remove_dir_all(&root); 4678 let runtime = RelayRuntime::open(runtime_config(&root, 8)).expect("runtime"); 4679 let rule = runtime.config().rate_limits().event().per_ip(); 4680 let saturated_peer_ip = IpAddr::V4(Ipv4Addr::new(127, 0, 0, 20)); 4681 let other_peer_ip = IpAddr::V4(Ipv4Addr::new(127, 0, 0, 21)); 4682 let key = TangleRateLimitKey::ip(TangleRateLimitScope::Event, saturated_peer_ip); 4683 for _ in 0..rule.max_hits() { 4684 runtime 4685 .rate_limiter() 4686 .record(key.clone(), rule, UnixTimestamp::new(1_714_124_433)); 4687 } 4688 let limited_event = 4689 tangle_v2_event(FixtureKey::Member, 1_714_124_433, 1, Vec::new(), "limited") 4690 .expect("limited event"); 4691 let rotated_event = 4692 tangle_v2_event(FixtureKey::Admin, 1_714_124_434, 2, Vec::new(), "rotated") 4693 .expect("rotated event"); 4694 let allowed_event = 4695 tangle_v2_event(FixtureKey::Owner, 1_714_124_435, 2, Vec::new(), "allowed") 4696 .expect("allowed event"); 4697 let handle = RelayRuntimeHandle::new(runtime); 4698 let mut auth = handle.auth_state().await.expect("auth"); 4699 4700 assert_eq!( 4701 handle 4702 .handle_protocol_client_message_with_rate_limit_context_for_test( 4703 ClientMessage::Event(limited_event.clone()), 4704 &mut auth, 4705 TangleClientRateLimitContext::new(Some(saturated_peer_ip), None), 4706 UnixTimestamp::new(1_714_124_433) 4707 ) 4708 .await 4709 .expect("event"), 4710 vec![RelayMessage::Ok { 4711 event_id: limited_event.id().clone(), 4712 accepted: false, 4713 message: "rate-limited: event ip rate limit exceeded until 1714124493".to_owned() 4714 }] 4715 ); 4716 assert_eq!( 4717 handle 4718 .handle_protocol_client_message_with_rate_limit_context_for_test( 4719 ClientMessage::Event(rotated_event.clone()), 4720 &mut auth, 4721 TangleClientRateLimitContext::new(Some(saturated_peer_ip), None), 4722 UnixTimestamp::new(1_714_124_433) 4723 ) 4724 .await 4725 .expect("event"), 4726 vec![RelayMessage::Ok { 4727 event_id: rotated_event.id().clone(), 4728 accepted: false, 4729 message: "rate-limited: event ip rate limit exceeded until 1714124493".to_owned() 4730 }] 4731 ); 4732 assert_eq!( 4733 handle 4734 .handle_protocol_client_message_with_rate_limit_context_for_test( 4735 ClientMessage::Event(allowed_event.clone()), 4736 &mut auth, 4737 TangleClientRateLimitContext::new(Some(other_peer_ip), None), 4738 UnixTimestamp::new(1_714_124_433) 4739 ) 4740 .await 4741 .expect("event"), 4742 vec![RelayMessage::Ok { 4743 event_id: allowed_event.id().clone(), 4744 accepted: true, 4745 message: String::new() 4746 }] 4747 ); 4748 assert_eq!(handle.metrics().rate_limit_rejections(), 2); 4749 assert_eq!(handle.metrics().event_rejections(), 2); 4750 assert_eq!(handle.metrics().event_admissions(), 1); 4751 4752 let _ = std::fs::remove_dir_all(root); 4753 } 4754 4755 #[tokio::test] 4756 async fn runtime_rate_limits_auth_pubkeys_before_authentication() { 4757 let root = temp_root("runtime-auth-pubkey-rate-limit"); 4758 let _ = std::fs::remove_dir_all(&root); 4759 let runtime = RelayRuntime::open(runtime_config(&root, 8)).expect("runtime"); 4760 let auth_event = 4761 tangle_v2_auth_event(FixtureKey::Member, "challenge-a", 120).expect("auth event"); 4762 let rule = runtime.config().rate_limits().auth().per_pubkey(); 4763 let key = TangleRateLimitKey::pubkey( 4764 TangleRateLimitScope::Auth, 4765 auth_event.unsigned().pubkey().clone(), 4766 ); 4767 for _ in 0..rule.max_hits() { 4768 runtime 4769 .rate_limiter() 4770 .record(key.clone(), rule, UnixTimestamp::new(120)); 4771 } 4772 let handle = RelayRuntimeHandle::new(runtime); 4773 let mut auth = handle.auth_state().await.expect("auth"); 4774 auth.issue_challenge("challenge-a", UnixTimestamp::new(100)) 4775 .expect("challenge"); 4776 4777 assert_eq!( 4778 handle 4779 .handle_protocol_client_message_for_test( 4780 ClientMessage::Auth(auth_event.clone()), 4781 &mut auth, 4782 UnixTimestamp::new(120) 4783 ) 4784 .await 4785 .expect("auth"), 4786 vec![RelayMessage::Ok { 4787 event_id: auth_event.id().clone(), 4788 accepted: false, 4789 message: "rate-limited: auth pubkey rate limit exceeded until 180".to_owned() 4790 }] 4791 ); 4792 assert!(auth.authenticated_pubkeys().is_empty()); 4793 4794 let _ = std::fs::remove_dir_all(root); 4795 } 4796 4797 #[tokio::test] 4798 async fn runtime_rate_limits_auth_peer_ips_before_authentication() { 4799 let root = temp_root("runtime-auth-ip-rate-limit"); 4800 let _ = std::fs::remove_dir_all(&root); 4801 let runtime = RelayRuntime::open(runtime_config(&root, 8)).expect("runtime"); 4802 let auth_event = 4803 tangle_v2_auth_event(FixtureKey::Member, "challenge-a", 120).expect("auth event"); 4804 let rule = runtime.config().rate_limits().auth().per_ip(); 4805 let peer_ip = IpAddr::V4(Ipv4Addr::new(127, 0, 0, 30)); 4806 let key = TangleRateLimitKey::ip(TangleRateLimitScope::Auth, peer_ip); 4807 for _ in 0..rule.max_hits() { 4808 runtime 4809 .rate_limiter() 4810 .record(key.clone(), rule, UnixTimestamp::new(120)); 4811 } 4812 let handle = RelayRuntimeHandle::new(runtime); 4813 let mut auth = handle.auth_state().await.expect("auth"); 4814 auth.issue_challenge("challenge-a", UnixTimestamp::new(100)) 4815 .expect("challenge"); 4816 4817 assert_eq!( 4818 handle 4819 .handle_protocol_client_message_with_rate_limit_context_for_test( 4820 ClientMessage::Auth(auth_event.clone()), 4821 &mut auth, 4822 TangleClientRateLimitContext::new(Some(peer_ip), None), 4823 UnixTimestamp::new(120) 4824 ) 4825 .await 4826 .expect("auth"), 4827 vec![RelayMessage::Ok { 4828 event_id: auth_event.id().clone(), 4829 accepted: false, 4830 message: "rate-limited: auth ip rate limit exceeded until 180".to_owned() 4831 }] 4832 ); 4833 assert!(auth.authenticated_pubkeys().is_empty()); 4834 assert_eq!(handle.metrics().rate_limit_rejections(), 1); 4835 assert_eq!(handle.metrics().auth_failures(), 1); 4836 4837 let _ = std::fs::remove_dir_all(root); 4838 } 4839 4840 #[tokio::test] 4841 async fn runtime_rate_limits_auth_failures() { 4842 let root = temp_root("runtime-auth-failure-rate-limit"); 4843 let _ = std::fs::remove_dir_all(&root); 4844 let runtime = RelayRuntime::open(runtime_config(&root, 8)).expect("runtime"); 4845 let auth_event = tangle_v2_event(FixtureKey::Member, 1_714_124_433, 22_242, Vec::new(), "") 4846 .expect("auth event"); 4847 let key = 4848 TangleRateLimitKey::auth_failure(None, Some(auth_event.unsigned().pubkey().clone())); 4849 let rule = runtime.config().rate_limits().auth().failures(); 4850 for _ in 0..rule.max_hits() { 4851 runtime 4852 .rate_limiter() 4853 .record(key.clone(), rule, UnixTimestamp::new(1_714_124_433)); 4854 } 4855 let handle = RelayRuntimeHandle::new(runtime); 4856 let mut auth = handle.auth_state().await.expect("auth"); 4857 4858 assert_eq!( 4859 handle 4860 .handle_protocol_client_message_for_test( 4861 ClientMessage::Auth(auth_event.clone()), 4862 &mut auth, 4863 UnixTimestamp::new(1_714_124_433) 4864 ) 4865 .await 4866 .expect("auth"), 4867 vec![RelayMessage::Ok { 4868 event_id: auth_event.id().clone(), 4869 accepted: false, 4870 message: "rate-limited: auth failure rate limit exceeded until 1714124733" 4871 .to_owned() 4872 }] 4873 ); 4874 4875 let _ = std::fs::remove_dir_all(root); 4876 } 4877 4878 #[tokio::test] 4879 async fn runtime_rate_limits_auth_failures_by_peer_ip() { 4880 let root = temp_root("runtime-auth-failure-ip-rate-limit"); 4881 let _ = std::fs::remove_dir_all(&root); 4882 let runtime = RelayRuntime::open(runtime_config(&root, 8)).expect("runtime"); 4883 let auth_event = tangle_v2_event(FixtureKey::Admin, 1_714_124_433, 22_242, Vec::new(), "") 4884 .expect("auth event"); 4885 let peer_ip = IpAddr::V4(Ipv4Addr::new(127, 0, 0, 31)); 4886 let key = TangleRateLimitKey::auth_failure(Some(peer_ip), None); 4887 let rule = runtime.config().rate_limits().auth().failures_per_ip(); 4888 for _ in 0..rule.max_hits() { 4889 runtime 4890 .rate_limiter() 4891 .record(key.clone(), rule, UnixTimestamp::new(1_714_124_433)); 4892 } 4893 let handle = RelayRuntimeHandle::new(runtime); 4894 let mut auth = handle.auth_state().await.expect("auth"); 4895 4896 assert_eq!( 4897 handle 4898 .handle_protocol_client_message_with_rate_limit_context_for_test( 4899 ClientMessage::Auth(auth_event.clone()), 4900 &mut auth, 4901 TangleClientRateLimitContext::new(Some(peer_ip), None), 4902 UnixTimestamp::new(1_714_124_433) 4903 ) 4904 .await 4905 .expect("auth"), 4906 vec![RelayMessage::Ok { 4907 event_id: auth_event.id().clone(), 4908 accepted: false, 4909 message: "rate-limited: auth failure ip rate limit exceeded until 1714124733" 4910 .to_owned() 4911 }] 4912 ); 4913 assert_eq!(handle.metrics().rate_limit_rejections(), 1); 4914 assert_eq!(handle.metrics().auth_failures(), 1); 4915 4916 let _ = std::fs::remove_dir_all(root); 4917 } 4918 4919 #[tokio::test] 4920 async fn runtime_preserves_chorus_auth_failure_rate_limit_parity() { 4921 let root = temp_root("runtime-chorus-auth-rate-limit-parity"); 4922 let _ = std::fs::remove_dir_all(&root); 4923 let runtime = RelayRuntime::open(runtime_config(&root, 8)).expect("runtime"); 4924 let pubkey_event = 4925 tangle_v2_event(FixtureKey::Member, 1_714_124_433, 22_242, Vec::new(), "") 4926 .expect("pubkey auth event"); 4927 let pubkey_rule = runtime.config().rate_limits().auth().failures(); 4928 let pubkey_key = 4929 TangleRateLimitKey::auth_failure(None, Some(pubkey_event.unsigned().pubkey().clone())); 4930 for _ in 0..pubkey_rule.max_hits() { 4931 runtime.rate_limiter().record( 4932 pubkey_key.clone(), 4933 pubkey_rule, 4934 UnixTimestamp::new(1_714_124_433), 4935 ); 4936 } 4937 let peer_ip = IpAddr::V4(Ipv4Addr::new(127, 0, 0, 41)); 4938 let peer_event = tangle_v2_event(FixtureKey::Admin, 1_714_124_434, 22_242, Vec::new(), "") 4939 .expect("peer auth event"); 4940 let peer_rule = runtime.config().rate_limits().auth().failures_per_ip(); 4941 let peer_key = TangleRateLimitKey::auth_failure(Some(peer_ip), None); 4942 for _ in 0..peer_rule.max_hits() { 4943 runtime.rate_limiter().record( 4944 peer_key.clone(), 4945 peer_rule, 4946 UnixTimestamp::new(1_714_124_434), 4947 ); 4948 } 4949 let handle = RelayRuntimeHandle::new(runtime); 4950 let mut auth = handle.auth_state().await.expect("auth"); 4951 4952 assert_eq!( 4953 handle 4954 .handle_protocol_client_message_for_test( 4955 ClientMessage::Auth(pubkey_event.clone()), 4956 &mut auth, 4957 UnixTimestamp::new(1_714_124_433) 4958 ) 4959 .await 4960 .expect("pubkey failure"), 4961 vec![RelayMessage::Ok { 4962 event_id: pubkey_event.id().clone(), 4963 accepted: false, 4964 message: "rate-limited: auth failure rate limit exceeded until 1714124733" 4965 .to_owned() 4966 }] 4967 ); 4968 assert_eq!( 4969 handle 4970 .handle_protocol_client_message_with_rate_limit_context_for_test( 4971 ClientMessage::Auth(peer_event.clone()), 4972 &mut auth, 4973 TangleClientRateLimitContext::new(Some(peer_ip), None), 4974 UnixTimestamp::new(1_714_124_434) 4975 ) 4976 .await 4977 .expect("peer failure"), 4978 vec![RelayMessage::Ok { 4979 event_id: peer_event.id().clone(), 4980 accepted: false, 4981 message: "rate-limited: auth failure ip rate limit exceeded until 1714124734" 4982 .to_owned() 4983 }] 4984 ); 4985 assert!(auth.authenticated_pubkeys().is_empty()); 4986 let snapshot = handle.metrics().snapshot(); 4987 assert_eq!(snapshot.client_messages(), 2); 4988 assert_eq!(snapshot.auth_messages(), 2); 4989 assert_eq!(snapshot.rate_limit_rejections(), 2); 4990 assert_eq!(handle.metrics().auth_successes(), 0); 4991 assert_eq!(handle.metrics().auth_failures(), 2); 4992 4993 let _ = std::fs::remove_dir_all(root); 4994 } 4995 4996 #[tokio::test] 4997 async fn runtime_rate_limits_group_writes_by_pubkey() { 4998 let root = temp_root("runtime-group-pubkey-rate-limit"); 4999 let _ = std::fs::remove_dir_all(&root); 5000 let runtime = RelayRuntime::open(runtime_config(&root, 8)).expect("runtime"); 5001 let event = tangle_v2_event( 5002 FixtureKey::Member, 5003 1_714_124_433, 5004 1, 5005 vec![Tag::from_parts("h", &["Farm"]).expect("h")], 5006 "limited", 5007 ) 5008 .expect("event"); 5009 let rule = runtime.config().rate_limits().group().write_per_pubkey(); 5010 let key = TangleRateLimitKey::pubkey( 5011 TangleRateLimitScope::GroupWrite, 5012 event.unsigned().pubkey().clone(), 5013 ); 5014 for _ in 0..rule.max_hits() { 5015 runtime 5016 .rate_limiter() 5017 .record(key.clone(), rule, UnixTimestamp::new(1_714_124_433)); 5018 } 5019 let handle = RelayRuntimeHandle::new(runtime); 5020 let mut auth = handle.auth_state().await.expect("auth"); 5021 5022 assert_eq!( 5023 handle 5024 .handle_protocol_client_message_for_test( 5025 ClientMessage::Event(event.clone()), 5026 &mut auth, 5027 UnixTimestamp::new(1_714_124_433) 5028 ) 5029 .await 5030 .expect("event"), 5031 vec![RelayMessage::Ok { 5032 event_id: event.id().clone(), 5033 accepted: false, 5034 message: "rate-limited: group pubkey rate limit exceeded until 1714124493" 5035 .to_owned() 5036 }] 5037 ); 5038 5039 let _ = std::fs::remove_dir_all(root); 5040 } 5041 5042 #[tokio::test] 5043 async fn runtime_rate_limits_group_writes_by_peer_ip() { 5044 let root = temp_root("runtime-group-ip-rate-limit"); 5045 let _ = std::fs::remove_dir_all(&root); 5046 let runtime = RelayRuntime::open(runtime_config(&root, 8)).expect("runtime"); 5047 let event = tangle_v2_event( 5048 FixtureKey::Member, 5049 1_714_124_433, 5050 1, 5051 vec![Tag::from_parts("h", &["Farm"]).expect("h")], 5052 "limited", 5053 ) 5054 .expect("event"); 5055 let rule = runtime.config().rate_limits().group().write_per_ip(); 5056 let peer_ip = IpAddr::V4(Ipv4Addr::new(127, 0, 0, 40)); 5057 let key = TangleRateLimitKey::ip(TangleRateLimitScope::GroupWrite, peer_ip); 5058 for _ in 0..rule.max_hits() { 5059 runtime 5060 .rate_limiter() 5061 .record(key.clone(), rule, UnixTimestamp::new(1_714_124_433)); 5062 } 5063 let handle = RelayRuntimeHandle::new(runtime); 5064 let mut auth = handle.auth_state().await.expect("auth"); 5065 5066 assert_eq!( 5067 handle 5068 .handle_protocol_client_message_with_rate_limit_context_for_test( 5069 ClientMessage::Event(event.clone()), 5070 &mut auth, 5071 TangleClientRateLimitContext::new(Some(peer_ip), None), 5072 UnixTimestamp::new(1_714_124_433) 5073 ) 5074 .await 5075 .expect("event"), 5076 vec![RelayMessage::Ok { 5077 event_id: event.id().clone(), 5078 accepted: false, 5079 message: "rate-limited: group ip rate limit exceeded until 1714124493".to_owned() 5080 }] 5081 ); 5082 assert_eq!(handle.metrics().rate_limit_rejections(), 1); 5083 assert_eq!(handle.metrics().event_rejections(), 1); 5084 assert_eq!(handle.metrics().group_write_denials(), 1); 5085 5086 let _ = std::fs::remove_dir_all(root); 5087 } 5088 5089 #[tokio::test] 5090 async fn runtime_rate_limits_group_writes_by_group_id() { 5091 let root = temp_root("runtime-group-write-rate-limit"); 5092 let _ = std::fs::remove_dir_all(&root); 5093 let runtime = RelayRuntime::open(runtime_config(&root, 8)).expect("runtime"); 5094 let group_id = GroupId::new("Farm").expect("group"); 5095 let event = tangle_v2_event( 5096 FixtureKey::Member, 5097 1_714_124_433, 5098 1, 5099 vec![Tag::from_parts("h", &[group_id.as_str()]).expect("h")], 5100 "limited", 5101 ) 5102 .expect("event"); 5103 let rule = runtime.config().rate_limits().group().write_per_group(); 5104 let key = TangleRateLimitKey::group(TangleRateLimitScope::GroupWrite, group_id); 5105 for _ in 0..rule.max_hits() { 5106 runtime 5107 .rate_limiter() 5108 .record(key.clone(), rule, UnixTimestamp::new(1_714_124_433)); 5109 } 5110 let handle = RelayRuntimeHandle::new(runtime); 5111 let mut auth = handle.auth_state().await.expect("auth"); 5112 5113 assert_eq!( 5114 handle 5115 .handle_protocol_client_message_for_test( 5116 ClientMessage::Event(event.clone()), 5117 &mut auth, 5118 UnixTimestamp::new(1_714_124_433) 5119 ) 5120 .await 5121 .expect("event"), 5122 vec![RelayMessage::Ok { 5123 event_id: event.id().clone(), 5124 accepted: false, 5125 message: "rate-limited: group write rate limit exceeded until 1714124493" 5126 .to_owned() 5127 }] 5128 ); 5129 5130 let _ = std::fs::remove_dir_all(root); 5131 } 5132 5133 #[tokio::test] 5134 async fn runtime_rate_limits_group_writes_by_kind() { 5135 let root = temp_root("runtime-group-kind-rate-limit"); 5136 let _ = std::fs::remove_dir_all(&root); 5137 let runtime = RelayRuntime::open(runtime_config(&root, 8)).expect("runtime"); 5138 let event = tangle_v2_event( 5139 FixtureKey::Member, 5140 1_714_124_433, 5141 1, 5142 vec![Tag::from_parts("h", &["Farm"]).expect("h")], 5143 "limited", 5144 ) 5145 .expect("event"); 5146 let rule = runtime.config().rate_limits().group().write_per_kind(); 5147 let key = 5148 TangleRateLimitKey::kind(TangleRateLimitScope::GroupWrite, event.unsigned().kind()); 5149 for _ in 0..rule.max_hits() { 5150 runtime 5151 .rate_limiter() 5152 .record(key.clone(), rule, UnixTimestamp::new(1_714_124_433)); 5153 } 5154 let handle = RelayRuntimeHandle::new(runtime); 5155 let mut auth = handle.auth_state().await.expect("auth"); 5156 5157 assert_eq!( 5158 handle 5159 .handle_protocol_client_message_for_test( 5160 ClientMessage::Event(event.clone()), 5161 &mut auth, 5162 UnixTimestamp::new(1_714_124_433) 5163 ) 5164 .await 5165 .expect("event"), 5166 vec![RelayMessage::Ok { 5167 event_id: event.id().clone(), 5168 accepted: false, 5169 message: "rate-limited: group kind rate limit exceeded until 1714124493".to_owned() 5170 }] 5171 ); 5172 5173 let _ = std::fs::remove_dir_all(root); 5174 } 5175 5176 #[tokio::test] 5177 async fn runtime_rate_limits_group_join_flows() { 5178 let root = temp_root("runtime-group-join-rate-limit"); 5179 let _ = std::fs::remove_dir_all(&root); 5180 let runtime = RelayRuntime::open(runtime_config(&root, 8)).expect("runtime"); 5181 let group_id = GroupId::new("Farm").expect("group"); 5182 let event = tangle_v2_event( 5183 FixtureKey::Member, 5184 1_714_124_433, 5185 KIND_GROUP_JOIN_REQUEST.into(), 5186 vec![Tag::from_parts("h", &[group_id.as_str()]).expect("h")], 5187 "", 5188 ) 5189 .expect("event"); 5190 let rule = runtime.config().rate_limits().group().join_flow(); 5191 let key = TangleRateLimitKey::join_flow(group_id, event.unsigned().pubkey().clone()); 5192 for _ in 0..rule.max_hits() { 5193 runtime 5194 .rate_limiter() 5195 .record(key.clone(), rule, UnixTimestamp::new(1_714_124_433)); 5196 } 5197 let handle = RelayRuntimeHandle::new(runtime); 5198 let mut auth = handle.auth_state().await.expect("auth"); 5199 5200 assert_eq!( 5201 handle 5202 .handle_protocol_client_message_for_test( 5203 ClientMessage::Event(event.clone()), 5204 &mut auth, 5205 UnixTimestamp::new(1_714_124_433) 5206 ) 5207 .await 5208 .expect("event"), 5209 vec![RelayMessage::Ok { 5210 event_id: event.id().clone(), 5211 accepted: false, 5212 message: "rate-limited: group join rate limit exceeded until 1714124733".to_owned() 5213 }] 5214 ); 5215 5216 let _ = std::fs::remove_dir_all(root); 5217 } 5218 5219 #[tokio::test] 5220 async fn runtime_rate_limits_group_join_flows_by_peer_ip() { 5221 let root = temp_root("runtime-group-join-ip-rate-limit"); 5222 let _ = std::fs::remove_dir_all(&root); 5223 let runtime = RelayRuntime::open(runtime_config(&root, 8)).expect("runtime"); 5224 let group_id = GroupId::new("Farm").expect("group"); 5225 let event = tangle_v2_event( 5226 FixtureKey::Member, 5227 1_714_124_433, 5228 KIND_GROUP_JOIN_REQUEST.into(), 5229 vec![Tag::from_parts("h", &[group_id.as_str()]).expect("h")], 5230 "", 5231 ) 5232 .expect("event"); 5233 let rule = runtime.config().rate_limits().group().join_flow_per_ip(); 5234 let peer_ip = IpAddr::V4(Ipv4Addr::new(127, 0, 0, 41)); 5235 let key = TangleRateLimitKey::join_flow_ip(group_id, peer_ip); 5236 for _ in 0..rule.max_hits() { 5237 runtime 5238 .rate_limiter() 5239 .record(key.clone(), rule, UnixTimestamp::new(1_714_124_433)); 5240 } 5241 let handle = RelayRuntimeHandle::new(runtime); 5242 let mut auth = handle.auth_state().await.expect("auth"); 5243 5244 assert_eq!( 5245 handle 5246 .handle_protocol_client_message_with_rate_limit_context_for_test( 5247 ClientMessage::Event(event.clone()), 5248 &mut auth, 5249 TangleClientRateLimitContext::new(Some(peer_ip), None), 5250 UnixTimestamp::new(1_714_124_433) 5251 ) 5252 .await 5253 .expect("event"), 5254 vec![RelayMessage::Ok { 5255 event_id: event.id().clone(), 5256 accepted: false, 5257 message: "rate-limited: group join ip rate limit exceeded until 1714124733" 5258 .to_owned() 5259 }] 5260 ); 5261 assert_eq!(handle.metrics().rate_limit_rejections(), 1); 5262 assert_eq!(handle.metrics().event_rejections(), 1); 5263 assert_eq!(handle.metrics().group_write_denials(), 1); 5264 5265 let _ = std::fs::remove_dir_all(root); 5266 } 5267 5268 #[tokio::test] 5269 async fn runtime_rate_limits_req_authenticated_pubkeys() { 5270 let root = temp_root("runtime-req-pubkey-rate-limit"); 5271 let _ = std::fs::remove_dir_all(&root); 5272 let runtime = RelayRuntime::open(runtime_config(&root, 8)).expect("runtime"); 5273 let rule = runtime.config().rate_limits().req().per_pubkey(); 5274 let handle = RelayRuntimeHandle::new(runtime); 5275 let mut auth = handle.auth_state().await.expect("auth"); 5276 auth.issue_challenge("challenge-a", UnixTimestamp::new(100)) 5277 .expect("challenge"); 5278 let auth_event = 5279 tangle_v2_auth_event(FixtureKey::Member, "challenge-a", 120).expect("auth event"); 5280 5281 assert_eq!( 5282 handle 5283 .handle_protocol_client_message_for_test( 5284 ClientMessage::Auth(auth_event.clone()), 5285 &mut auth, 5286 UnixTimestamp::new(120) 5287 ) 5288 .await 5289 .expect("auth"), 5290 vec![RelayMessage::Ok { 5291 event_id: auth_event.id().clone(), 5292 accepted: true, 5293 message: String::new() 5294 }] 5295 ); 5296 let key = 5297 TangleRateLimitKey::pubkey(TangleRateLimitScope::Req, FixtureKey::Member.public_key()); 5298 let limiter = handle.rate_limiter().await; 5299 for _ in 0..rule.max_hits() { 5300 limiter.record(key.clone(), rule, UnixTimestamp::new(120)); 5301 } 5302 let subscription_id = SubscriptionId::new("limited-req-pubkey").expect("subscription"); 5303 let filters = vec![filter_from_value(&json!({"limit": 1})).expect("filter")]; 5304 5305 assert_eq!( 5306 handle 5307 .handle_protocol_client_message_for_test( 5308 ClientMessage::Req { 5309 subscription_id: subscription_id.clone(), 5310 filters 5311 }, 5312 &mut auth, 5313 UnixTimestamp::new(120) 5314 ) 5315 .await 5316 .expect("req"), 5317 vec![RelayMessage::Closed { 5318 subscription_id, 5319 message: "rate-limited: req pubkey rate limit exceeded until 180".to_owned() 5320 }] 5321 ); 5322 5323 let _ = std::fs::remove_dir_all(root); 5324 } 5325 5326 #[tokio::test] 5327 async fn runtime_rate_limits_req_connections() { 5328 let root = temp_root("runtime-req-connection-rate-limit"); 5329 let _ = std::fs::remove_dir_all(&root); 5330 let runtime = RelayRuntime::open(runtime_config(&root, 8)).expect("runtime"); 5331 let rule = runtime.config().rate_limits().req().per_connection(); 5332 let key = TangleRateLimitKey::connection(TangleRateLimitScope::Req, 77); 5333 for _ in 0..rule.max_hits() { 5334 runtime 5335 .rate_limiter() 5336 .record(key.clone(), rule, UnixTimestamp::new(1_714_124_433)); 5337 } 5338 let handle = RelayRuntimeHandle::new(runtime); 5339 let mut auth = handle.auth_state().await.expect("auth"); 5340 let subscription_id = SubscriptionId::new("limited-req-connection").expect("subscription"); 5341 let filters = vec![filter_from_value(&json!({"kinds": [1], "limit": 1})).expect("filter")]; 5342 5343 assert_eq!( 5344 handle 5345 .handle_protocol_client_message_with_rate_limit_context_for_test( 5346 ClientMessage::Req { 5347 subscription_id: subscription_id.clone(), 5348 filters 5349 }, 5350 &mut auth, 5351 TangleClientRateLimitContext::new(None, Some(77)), 5352 UnixTimestamp::new(1_714_124_433) 5353 ) 5354 .await 5355 .expect("req"), 5356 vec![RelayMessage::Closed { 5357 subscription_id, 5358 message: "rate-limited: req connection rate limit exceeded until 1714124493" 5359 .to_owned() 5360 }] 5361 ); 5362 5363 let _ = std::fs::remove_dir_all(root); 5364 } 5365 5366 #[tokio::test] 5367 async fn runtime_rate_limits_req_filter_groups() { 5368 let root = temp_root("runtime-req-group-rate-limit"); 5369 let _ = std::fs::remove_dir_all(&root); 5370 let runtime = RelayRuntime::open(runtime_config(&root, 8)).expect("runtime"); 5371 let group_id = GroupId::new("Farm").expect("group"); 5372 let rule = runtime.config().rate_limits().req().per_group(); 5373 let key = TangleRateLimitKey::group(TangleRateLimitScope::Req, group_id); 5374 for _ in 0..rule.max_hits() { 5375 runtime 5376 .rate_limiter() 5377 .record(key.clone(), rule, UnixTimestamp::new(1_714_124_433)); 5378 } 5379 let handle = RelayRuntimeHandle::new(runtime); 5380 let mut auth = handle.auth_state().await.expect("auth"); 5381 let subscription_id = SubscriptionId::new("limited-req-group").expect("subscription"); 5382 let filters = 5383 vec![filter_from_value(&json!({"#h": ["Farm"], "limit": 1})).expect("filter")]; 5384 5385 assert_eq!( 5386 handle 5387 .handle_protocol_client_message_for_test( 5388 ClientMessage::Req { 5389 subscription_id: subscription_id.clone(), 5390 filters 5391 }, 5392 &mut auth, 5393 UnixTimestamp::new(1_714_124_433) 5394 ) 5395 .await 5396 .expect("req"), 5397 vec![RelayMessage::Closed { 5398 subscription_id, 5399 message: "rate-limited: req group rate limit exceeded until 1714124493".to_owned() 5400 }] 5401 ); 5402 5403 let _ = std::fs::remove_dir_all(root); 5404 } 5405 5406 #[test] 5407 fn query_classifier_identifies_broad_count_shapes() { 5408 let classifier = TangleQueryClassifier::new(runtime_relay_limits(8)); 5409 let empty_filter = pocket_filter(json!({})); 5410 let tag_only_filter = pocket_filter(json!({"#t": ["market"], "limit": 1})); 5411 let kind_only_filter = pocket_filter(json!({"kinds": [1], "limit": 1})); 5412 let high_limit_filter = pocket_filter(json!({"kinds": [1], "#h": ["Farm"], "limit": 500})); 5413 let broad_time_filter = pocket_filter(json!({ 5414 "kinds": [1], 5415 "since": 1, 5416 "until": BROAD_QUERY_TIME_WINDOW_SECONDS + 2, 5417 "limit": 1 5418 })); 5419 let bounded_group_filter = pocket_filter(json!({"kinds": [1], "#h": ["Farm"], "limit": 1})); 5420 let bounded_time_filter = pocket_filter(json!({ 5421 "kinds": [1], 5422 "since": 1, 5423 "until": BROAD_QUERY_TIME_WINDOW_SECONDS, 5424 "limit": 1 5425 })); 5426 let hll_reaction_filter = pocket_filter(json!({"kinds": [7], "#e": ["a".repeat(64)]})); 5427 5428 assert_eq!( 5429 classifier.classify_pocket_count(&[]), 5430 TangleQueryClassification::Broad(TangleBroadQueryReason::EmptyFilters) 5431 ); 5432 assert_eq!( 5433 classifier.classify_pocket_count(&[empty_filter]), 5434 TangleQueryClassification::Broad(TangleBroadQueryReason::MissingPrimaryConstraint) 5435 ); 5436 assert_eq!( 5437 classifier.classify_pocket_count(&[tag_only_filter]), 5438 TangleQueryClassification::Broad(TangleBroadQueryReason::MissingPrimaryConstraint) 5439 ); 5440 assert_eq!( 5441 classifier.classify_pocket_count(&[kind_only_filter]), 5442 TangleQueryClassification::Broad(TangleBroadQueryReason::MissingBoundedSelector) 5443 ); 5444 assert_eq!( 5445 classifier.classify_pocket_count(&[high_limit_filter]), 5446 TangleQueryClassification::Broad(TangleBroadQueryReason::HighLimit) 5447 ); 5448 assert_eq!( 5449 classifier.classify_pocket_count(&[broad_time_filter]), 5450 TangleQueryClassification::Broad(TangleBroadQueryReason::BroadTimeWindow) 5451 ); 5452 assert_eq!( 5453 classifier.classify_pocket_count(&[bounded_group_filter]), 5454 TangleQueryClassification::Bounded 5455 ); 5456 assert_eq!( 5457 classifier.classify_pocket_count(&[bounded_time_filter]), 5458 TangleQueryClassification::Bounded 5459 ); 5460 assert_eq!( 5461 classifier.classify_pocket_count(&[hll_reaction_filter]), 5462 TangleQueryClassification::Bounded 5463 ); 5464 } 5465 5466 #[tokio::test] 5467 async fn runtime_count_hll_accepts_public_pocket_selector() { 5468 let root = temp_root("runtime-count-hll"); 5469 let _ = std::fs::remove_dir_all(&root); 5470 let handle = 5471 RelayRuntimeHandle::new(RelayRuntime::open(runtime_config(&root, 8)).expect("runtime")); 5472 let mut auth = handle.auth_state().await.expect("auth"); 5473 let target = "c".repeat(64); 5474 let tags = PocketOwnedTags::new(&[["e", target.as_str()]]).expect("tags"); 5475 let first = signed_pocket_event(12, 1_714_124_433, 7, &tags, b"first reaction"); 5476 let second = signed_pocket_event(11, 1_714_124_434, 7, &tags, b"second reaction"); 5477 5478 assert_accepted_pocket_reply( 5479 runtime_pocket_event_reply(&handle, &first, &mut auth), 5480 &first, 5481 ); 5482 assert_accepted_pocket_reply( 5483 runtime_pocket_event_reply(&handle, &second, &mut auth), 5484 &second, 5485 ); 5486 5487 let subscription_id = SubscriptionId::new("count-hll-runtime").expect("subscription"); 5488 let replies = handle 5489 .handle_protocol_client_message_for_test( 5490 ClientMessage::Count { 5491 subscription_id: subscription_id.clone(), 5492 filters: vec![ 5493 filter_from_value(&json!({"kinds":[7],"#e":[target]})).expect("filter"), 5494 ], 5495 }, 5496 &mut auth, 5497 UnixTimestamp::new(1_714_124_437), 5498 ) 5499 .await 5500 .expect("count"); 5501 let [ 5502 RelayMessage::Count { 5503 subscription_id: actual, 5504 count, 5505 hll: Some(hll), 5506 }, 5507 ] = replies.as_slice() 5508 else { 5509 panic!("count hll expected: {replies:?}") 5510 }; 5511 5512 assert_eq!(actual, &subscription_id); 5513 assert_eq!(*count, 2); 5514 assert_eq!(hll.len(), 512); 5515 assert_ne!(hll, &"00".repeat(256)); 5516 5517 let _ = std::fs::remove_dir_all(root); 5518 } 5519 5520 #[test] 5521 fn runtime_count_source_stays_exact() { 5522 let sources = [ 5523 include_str!("runtime.rs"), 5524 include_str!("relay/core.rs"), 5525 include_str!("../../tangle_protocol/src/lib.rs"), 5526 ]; 5527 let forbidden = [ 5528 concat!("approximate", "_count"), 5529 concat!("approx", "_count"), 5530 concat!("estimated", "_count"), 5531 concat!("count", "_estimate"), 5532 concat!("private", "_count", "_estimate"), 5533 ]; 5534 5535 for source in sources { 5536 for needle in forbidden { 5537 assert!(!source.contains(needle)); 5538 } 5539 } 5540 } 5541 5542 #[tokio::test] 5543 async fn runtime_rate_limits_count_peer_ips() { 5544 let root = temp_root("runtime-count-ip-rate-limit"); 5545 let _ = std::fs::remove_dir_all(&root); 5546 let runtime = RelayRuntime::open(runtime_config(&root, 8)).expect("runtime"); 5547 let rule = runtime.config().rate_limits().count().per_ip(); 5548 let peer_ip = IpAddr::V4(Ipv4Addr::new(127, 0, 0, 9)); 5549 let key = TangleRateLimitKey::ip(TangleRateLimitScope::Count, peer_ip); 5550 for _ in 0..rule.max_hits() { 5551 runtime 5552 .rate_limiter() 5553 .record(key.clone(), rule, UnixTimestamp::new(1_714_124_433)); 5554 } 5555 let handle = RelayRuntimeHandle::new(runtime); 5556 let mut auth = handle.auth_state().await.expect("auth"); 5557 let subscription_id = SubscriptionId::new("limited-count-ip").expect("subscription"); 5558 let filters = vec![ 5559 filter_from_value(&json!({"kinds": [1], "#h": ["Farm"], "limit": 1})).expect("filter"), 5560 ]; 5561 5562 assert_eq!( 5563 handle 5564 .handle_protocol_client_message_with_rate_limit_context_for_test( 5565 ClientMessage::Count { 5566 subscription_id: subscription_id.clone(), 5567 filters 5568 }, 5569 &mut auth, 5570 TangleClientRateLimitContext::new(Some(peer_ip), None), 5571 UnixTimestamp::new(1_714_124_433) 5572 ) 5573 .await 5574 .expect("count"), 5575 vec![RelayMessage::Closed { 5576 subscription_id, 5577 message: "rate-limited: count ip rate limit exceeded until 1714124493".to_owned() 5578 }] 5579 ); 5580 5581 let _ = std::fs::remove_dir_all(root); 5582 } 5583 5584 #[tokio::test] 5585 async fn runtime_rejects_search_req_and_count_as_unsupported() { 5586 let root = temp_root("runtime-search-unsupported"); 5587 let _ = std::fs::remove_dir_all(&root); 5588 let handle = 5589 RelayRuntimeHandle::new(RelayRuntime::open(runtime_config(&root, 8)).expect("runtime")); 5590 let mut auth = handle.auth_state().await.expect("auth"); 5591 let req_id = SubscriptionId::new("search-req").expect("req"); 5592 let count_id = SubscriptionId::new("search-count").expect("count"); 5593 let search = 5594 filter_from_value(&json!({"search": "fresh carrots", "limit": 1})).expect("filter"); 5595 5596 assert_eq!( 5597 handle 5598 .handle_protocol_client_message_for_test( 5599 ClientMessage::Req { 5600 subscription_id: req_id.clone(), 5601 filters: vec![search.clone()] 5602 }, 5603 &mut auth, 5604 UnixTimestamp::new(1_714_124_433) 5605 ) 5606 .await 5607 .expect("req"), 5608 vec![RelayMessage::Closed { 5609 subscription_id: req_id, 5610 message: "unsupported: search filters are not supported".to_owned() 5611 }] 5612 ); 5613 assert_eq!( 5614 handle 5615 .handle_protocol_client_message_for_test( 5616 ClientMessage::Count { 5617 subscription_id: count_id.clone(), 5618 filters: vec![search] 5619 }, 5620 &mut auth, 5621 UnixTimestamp::new(1_714_124_434) 5622 ) 5623 .await 5624 .expect("count"), 5625 vec![RelayMessage::Closed { 5626 subscription_id: count_id, 5627 message: "unsupported: search filters are not supported".to_owned() 5628 }] 5629 ); 5630 5631 let _ = std::fs::remove_dir_all(root); 5632 } 5633 5634 #[tokio::test] 5635 async fn runtime_rate_limits_count_filter_kinds() { 5636 let root = temp_root("runtime-count-kind-rate-limit"); 5637 let _ = std::fs::remove_dir_all(&root); 5638 let runtime = RelayRuntime::open(runtime_config(&root, 8)).expect("runtime"); 5639 let kind = Kind::new(1).expect("kind"); 5640 let rule = runtime.config().rate_limits().count().per_kind(); 5641 let key = TangleRateLimitKey::kind(TangleRateLimitScope::Count, kind); 5642 for _ in 0..rule.max_hits() { 5643 runtime 5644 .rate_limiter() 5645 .record(key.clone(), rule, UnixTimestamp::new(1_714_124_433)); 5646 } 5647 let handle = RelayRuntimeHandle::new(runtime); 5648 let mut auth = handle.auth_state().await.expect("auth"); 5649 let subscription_id = SubscriptionId::new("limited-count-kind").expect("subscription"); 5650 let filters = vec![ 5651 filter_from_value(&json!({"kinds": [1], "#h": ["Farm"], "limit": 1})).expect("filter"), 5652 ]; 5653 5654 assert_eq!( 5655 handle 5656 .handle_protocol_client_message_for_test( 5657 ClientMessage::Count { 5658 subscription_id: subscription_id.clone(), 5659 filters 5660 }, 5661 &mut auth, 5662 UnixTimestamp::new(1_714_124_433) 5663 ) 5664 .await 5665 .expect("count"), 5666 vec![RelayMessage::Closed { 5667 subscription_id, 5668 message: "rate-limited: count kind rate limit exceeded until 1714124493".to_owned() 5669 }] 5670 ); 5671 5672 let _ = std::fs::remove_dir_all(root); 5673 } 5674 5675 #[tokio::test] 5676 async fn runtime_refuses_broad_count_queries_before_rate_limits() { 5677 let root = temp_root("runtime-count-broad-refusal"); 5678 let _ = std::fs::remove_dir_all(&root); 5679 let runtime = RelayRuntime::open(runtime_config(&root, 8)).expect("runtime"); 5680 let rule = runtime.config().rate_limits().count().broad(); 5681 let key = TangleRateLimitKey::query_class( 5682 TangleRateLimitScope::Count, 5683 TangleRateLimitQueryClass::Broad, 5684 ); 5685 for _ in 0..rule.max_hits() { 5686 runtime 5687 .rate_limiter() 5688 .record(key.clone(), rule, UnixTimestamp::new(1_714_124_433)); 5689 } 5690 let handle = RelayRuntimeHandle::new(runtime); 5691 let mut auth = handle.auth_state().await.expect("auth"); 5692 let subscription_id = SubscriptionId::new("limited-count-broad").expect("subscription"); 5693 let filters = vec![filter_from_value(&json!({"limit": 1})).expect("filter")]; 5694 5695 assert_eq!( 5696 handle 5697 .handle_protocol_client_message_for_test( 5698 ClientMessage::Count { 5699 subscription_id: subscription_id.clone(), 5700 filters 5701 }, 5702 &mut auth, 5703 UnixTimestamp::new(1_714_124_433) 5704 ) 5705 .await 5706 .expect("count"), 5707 vec![RelayMessage::Closed { 5708 subscription_id, 5709 message: "restricted: count filters are too broad or expensive".to_owned() 5710 }] 5711 ); 5712 assert_eq!(handle.metrics().count_refusals(), 1); 5713 assert_eq!(handle.metrics().broad_query_rejections(), 1); 5714 5715 let _ = std::fs::remove_dir_all(root); 5716 } 5717 5718 #[tokio::test] 5719 async fn runtime_refuses_expensive_count_queries_deterministically() { 5720 let root = temp_root("runtime-count-expensive-refusal"); 5721 let _ = std::fs::remove_dir_all(&root); 5722 let handle = 5723 RelayRuntimeHandle::new(RelayRuntime::open(runtime_config(&root, 8)).expect("runtime")); 5724 let mut auth = handle.auth_state().await.expect("auth"); 5725 let cases = [ 5726 ("missing-selector", json!({"kinds": [1], "limit": 1})), 5727 ( 5728 "high-limit", 5729 json!({"kinds": [1], "#h": ["Farm"], "limit": 500}), 5730 ), 5731 ( 5732 "broad-window", 5733 json!({ 5734 "kinds": [1], 5735 "since": 1, 5736 "until": BROAD_QUERY_TIME_WINDOW_SECONDS + 2, 5737 "limit": 1 5738 }), 5739 ), 5740 ]; 5741 5742 for (name, value) in cases { 5743 let subscription_id = SubscriptionId::new(name).expect("subscription"); 5744 let filters = vec![filter_from_value(&value).expect("filter")]; 5745 5746 assert_eq!( 5747 handle 5748 .handle_protocol_client_message_for_test( 5749 ClientMessage::Count { 5750 subscription_id: subscription_id.clone(), 5751 filters 5752 }, 5753 &mut auth, 5754 UnixTimestamp::new(1_714_124_433) 5755 ) 5756 .await 5757 .expect("count"), 5758 vec![RelayMessage::Closed { 5759 subscription_id, 5760 message: "restricted: count filters are too broad or expensive".to_owned() 5761 }] 5762 ); 5763 } 5764 assert_eq!(handle.metrics().count_refusals(), 3); 5765 assert_eq!(handle.metrics().broad_query_rejections(), 3); 5766 5767 let _ = std::fs::remove_dir_all(root); 5768 } 5769 5770 #[tokio::test] 5771 async fn runtime_publishes_generated_group_event_offsets_for_live_fanout() { 5772 let root = temp_root("runtime-generated-offset-fanout"); 5773 let _ = std::fs::remove_dir_all(&root); 5774 let handle = 5775 RelayRuntimeHandle::new(RelayRuntime::open(runtime_config(&root, 8)).expect("runtime")); 5776 let mut offsets = handle.subscribe_events().await; 5777 let mut auth = handle.auth_state().await.expect("auth"); 5778 auth.issue_challenge("challenge-a", UnixTimestamp::new(100)) 5779 .expect("challenge"); 5780 let auth_event = 5781 tangle_v2_auth_event(FixtureKey::Owner, "challenge-a", 120).expect("auth event"); 5782 let create = tangle_v2_group_create_event(FixtureKey::Owner, "RuntimeFarm", 121, &[]) 5783 .expect("create"); 5784 let mut subscriptions = LiveSubscriptionSet::new(8, 64).expect("subscriptions"); 5785 let subscription_id = SubscriptionId::new("generated-offsets").expect("subscription"); 5786 subscriptions 5787 .subscribe( 5788 subscription_id.clone(), 5789 vec![pocket_filter(json!({ 5790 "kinds":[KIND_GROUP_METADATA, KIND_GROUP_ADMINS, KIND_GROUP_MEMBERS], 5791 "#d":["RuntimeFarm"] 5792 }))], 5793 ) 5794 .expect("subscribe"); 5795 5796 assert_eq!( 5797 handle 5798 .handle_protocol_client_message_for_test( 5799 ClientMessage::Auth(auth_event.clone()), 5800 &mut auth, 5801 UnixTimestamp::new(120) 5802 ) 5803 .await 5804 .expect("auth"), 5805 vec![RelayMessage::Ok { 5806 event_id: auth_event.id().clone(), 5807 accepted: true, 5808 message: String::new() 5809 }] 5810 ); 5811 assert_eq!( 5812 handle 5813 .handle_protocol_client_message_for_test( 5814 ClientMessage::Event(create.clone()), 5815 &mut auth, 5816 UnixTimestamp::new(121) 5817 ) 5818 .await 5819 .expect("create"), 5820 vec![RelayMessage::Ok { 5821 event_id: create.id().clone(), 5822 accepted: true, 5823 message: String::new() 5824 }] 5825 ); 5826 let source_offset = offsets.try_recv().expect("source offset"); 5827 let generated_offsets = [ 5828 offsets.try_recv().expect("first generated offset"), 5829 offsets.try_recv().expect("second generated offset"), 5830 ]; 5831 assert!(source_offset < generated_offsets[0]); 5832 assert!(generated_offsets[0] < generated_offsets[1]); 5833 let put_member = 5834 tangle_v2_put_user_event(FixtureKey::Owner, "RuntimeFarm", FixtureKey::Member, 122) 5835 .expect("put member"); 5836 assert_eq!( 5837 handle 5838 .handle_protocol_client_message_for_test( 5839 ClientMessage::Event(put_member.clone()), 5840 &mut auth, 5841 UnixTimestamp::new(122) 5842 ) 5843 .await 5844 .expect("put member"), 5845 vec![RelayMessage::Ok { 5846 event_id: put_member.id().clone(), 5847 accepted: true, 5848 message: String::new() 5849 }] 5850 ); 5851 let put_source_offset = offsets.try_recv().expect("put source offset"); 5852 let member_generated_offset = offsets.try_recv().expect("member generated offset"); 5853 assert!(generated_offsets[1] < put_source_offset); 5854 assert!(put_source_offset < member_generated_offset); 5855 let generated_offsets = [ 5856 generated_offsets[0], 5857 generated_offsets[1], 5858 member_generated_offset, 5859 ]; 5860 let mut generated_kinds = BTreeSet::new(); 5861 for offset in generated_offsets { 5862 let messages = handle 5863 .fanout_event_offset_with_projection_context( 5864 offset, 5865 &mut subscriptions, 5866 &auth, 5867 &RelayProjectionContext::default(), 5868 ) 5869 .await 5870 .expect("fanout"); 5871 assert!(matches!( 5872 messages.as_slice(), 5873 [RuntimeRelayMessage::Event { 5874 subscription_id: delivered, 5875 event 5876 }] if delivered == &subscription_id 5877 && generated_kinds.insert(u32::from(event.kind().as_u16())) 5878 )); 5879 } 5880 assert_eq!( 5881 generated_kinds, 5882 BTreeSet::from([KIND_GROUP_METADATA, KIND_GROUP_ADMINS, KIND_GROUP_MEMBERS]) 5883 ); 5884 assert_eq!(handle.metrics().outbox_replayed_events(), 3); 5885 assert_eq!(handle.metrics().outbox_pending_events(), 0); 5886 assert_eq!(handle.metrics().event_bus_published_offsets(), 5); 5887 assert_eq!( 5888 offsets.try_recv().expect_err("only source plus generated"), 5889 TangleEventReceiveError::Empty 5890 ); 5891 5892 let _ = std::fs::remove_dir_all(root); 5893 } 5894 5895 #[tokio::test] 5896 async fn runtime_group_concurrency_duplicate_create_accepts_one_projection() { 5897 let root = temp_root("runtime-group-concurrency-duplicate-create"); 5898 let _ = std::fs::remove_dir_all(&root); 5899 let handle = RelayRuntimeHandle::new( 5900 RelayRuntime::open(runtime_config(&root, 32)).expect("runtime"), 5901 ); 5902 let mut offsets = handle.subscribe_events().await; 5903 let owner_auth = 5904 authenticated_runtime_state(&handle, FixtureKey::Owner, "owner-create", 1_714_126_100) 5905 .await; 5906 let first = 5907 tangle_v2_group_create_event(FixtureKey::Owner, "RaceCreate", 1_714_126_101, &[]) 5908 .expect("first create"); 5909 let second = 5910 tangle_v2_group_create_event(FixtureKey::Owner, "RaceCreate", 1_714_126_102, &[]) 5911 .expect("second create"); 5912 let first_task = { 5913 let handle = handle.clone(); 5914 let mut auth = owner_auth.clone(); 5915 let event = first.clone(); 5916 tokio::spawn(async move { 5917 runtime_event_reply(&handle, event, &mut auth, 1_714_126_101).await 5918 }) 5919 }; 5920 let second_task = { 5921 let handle = handle.clone(); 5922 let mut auth = owner_auth.clone(); 5923 let event = second.clone(); 5924 tokio::spawn(async move { 5925 runtime_event_reply(&handle, event, &mut auth, 1_714_126_102).await 5926 }) 5927 }; 5928 let replies = tokio::time::timeout(Duration::from_secs(3), async { 5929 vec![ 5930 first_task.await.expect("first task"), 5931 second_task.await.expect("second task"), 5932 ] 5933 }) 5934 .await 5935 .expect("duplicate create race"); 5936 5937 assert_eq!(accepted_count(&replies), 1); 5938 assert_eq!( 5939 rejected_messages(&replies), 5940 vec!["invalid: group already exists".to_owned()] 5941 ); 5942 assert_eq!(drain_offsets(&mut offsets, 3).await.len(), 3); 5943 assert_eq!( 5944 offsets 5945 .try_recv() 5946 .expect_err("one create source plus generated"), 5947 TangleEventReceiveError::Empty 5948 ); 5949 let mut auth = owner_auth.clone(); 5950 assert_eq!( 5951 runtime_group_count( 5952 &handle, 5953 "duplicate-create-count", 5954 "RaceCreate", 5955 KIND_GROUP_METADATA, 5956 "d", 5957 &mut auth, 5958 1_714_126_103, 5959 ) 5960 .await, 5961 1 5962 ); 5963 assert_live_projection_matches_rebuild(&handle, "RaceCreate"); 5964 5965 let _ = std::fs::remove_dir_all(root); 5966 } 5967 5968 #[tokio::test] 5969 async fn runtime_group_concurrency_duplicate_join_accepts_one_membership() { 5970 let root = temp_root("runtime-group-concurrency-duplicate-join"); 5971 let _ = std::fs::remove_dir_all(&root); 5972 let handle = RelayRuntimeHandle::new( 5973 RelayRuntime::open(runtime_config_with_public_join(&root, 32)).expect("runtime"), 5974 ); 5975 let mut offsets = handle.subscribe_events().await; 5976 let mut owner_auth = 5977 authenticated_runtime_state(&handle, FixtureKey::Owner, "owner-join", 1_714_126_200) 5978 .await; 5979 let member_auth = 5980 authenticated_runtime_state(&handle, FixtureKey::Member, "member-join", 1_714_126_201) 5981 .await; 5982 let create = 5983 tangle_v2_group_create_event(FixtureKey::Owner, "RaceJoin", 1_714_126_202, &[]) 5984 .expect("create"); 5985 assert_accepted_reply( 5986 runtime_event_reply(&handle, create.clone(), &mut owner_auth, 1_714_126_202).await, 5987 &create, 5988 ); 5989 assert_eq!(drain_offsets(&mut offsets, 3).await.len(), 3); 5990 let join_a = 5991 tangle_v2_join_event(FixtureKey::Member, "RaceJoin", 1_714_126_203).expect("join a"); 5992 let join_b = 5993 tangle_v2_join_event(FixtureKey::Member, "RaceJoin", 1_714_126_204).expect("join b"); 5994 let first_task = { 5995 let handle = handle.clone(); 5996 let mut auth = member_auth.clone(); 5997 let event = join_a.clone(); 5998 tokio::spawn(async move { 5999 runtime_event_reply(&handle, event, &mut auth, 1_714_126_203).await 6000 }) 6001 }; 6002 let second_task = { 6003 let handle = handle.clone(); 6004 let mut auth = member_auth.clone(); 6005 let event = join_b.clone(); 6006 tokio::spawn(async move { 6007 runtime_event_reply(&handle, event, &mut auth, 1_714_126_204).await 6008 }) 6009 }; 6010 let replies = tokio::time::timeout(Duration::from_secs(3), async { 6011 vec![ 6012 first_task.await.expect("first task"), 6013 second_task.await.expect("second task"), 6014 ] 6015 }) 6016 .await 6017 .expect("duplicate join race"); 6018 6019 assert_eq!(accepted_count(&replies), 1); 6020 assert_eq!( 6021 rejected_messages(&replies), 6022 vec!["duplicate: group member already exists".to_owned()] 6023 ); 6024 assert_eq!(drain_offsets(&mut offsets, 2).await.len(), 2); 6025 assert_runtime_member_status( 6026 &handle, 6027 "RaceJoin", 6028 &FixtureKey::Member.public_key(), 6029 MemberStatus::Member, 6030 ); 6031 assert_live_projection_matches_rebuild(&handle, "RaceJoin"); 6032 6033 let _ = std::fs::remove_dir_all(root); 6034 } 6035 6036 #[tokio::test] 6037 async fn runtime_group_concurrency_join_and_leave_match_rebuild() { 6038 let root = temp_root("runtime-group-concurrency-join-leave"); 6039 let _ = std::fs::remove_dir_all(&root); 6040 let handle = RelayRuntimeHandle::new( 6041 RelayRuntime::open(runtime_config_with_public_join(&root, 32)).expect("runtime"), 6042 ); 6043 let mut owner_auth = authenticated_runtime_state( 6044 &handle, 6045 FixtureKey::Owner, 6046 "owner-join-leave", 6047 1_714_126_300, 6048 ) 6049 .await; 6050 let member_auth = authenticated_runtime_state( 6051 &handle, 6052 FixtureKey::Member, 6053 "member-join-leave", 6054 1_714_126_301, 6055 ) 6056 .await; 6057 let create = 6058 tangle_v2_group_create_event(FixtureKey::Owner, "RaceJoinLeave", 1_714_126_302, &[]) 6059 .expect("create"); 6060 let put_member = tangle_v2_put_user_event( 6061 FixtureKey::Owner, 6062 "RaceJoinLeave", 6063 FixtureKey::Member, 6064 1_714_126_303, 6065 ) 6066 .expect("put member"); 6067 assert_accepted_reply( 6068 runtime_event_reply(&handle, create.clone(), &mut owner_auth, 1_714_126_302).await, 6069 &create, 6070 ); 6071 assert_accepted_reply( 6072 runtime_event_reply(&handle, put_member.clone(), &mut owner_auth, 1_714_126_303).await, 6073 &put_member, 6074 ); 6075 let leave = tangle_v2_leave_event(FixtureKey::Member, "RaceJoinLeave", 1_714_126_304) 6076 .expect("leave"); 6077 let join = 6078 tangle_v2_join_event(FixtureKey::Member, "RaceJoinLeave", 1_714_126_305).expect("join"); 6079 let leave_task = { 6080 let handle = handle.clone(); 6081 let mut auth = member_auth.clone(); 6082 let event = leave.clone(); 6083 tokio::spawn(async move { 6084 runtime_event_reply(&handle, event, &mut auth, 1_714_126_304).await 6085 }) 6086 }; 6087 let join_task = { 6088 let handle = handle.clone(); 6089 let mut auth = member_auth.clone(); 6090 let event = join.clone(); 6091 tokio::spawn(async move { 6092 runtime_event_reply(&handle, event, &mut auth, 1_714_126_305).await 6093 }) 6094 }; 6095 let replies = tokio::time::timeout(Duration::from_secs(3), async { 6096 vec![ 6097 leave_task.await.expect("leave task"), 6098 join_task.await.expect("join task"), 6099 ] 6100 }) 6101 .await 6102 .expect("join leave race"); 6103 let join_accepted = reply_is_accepted(&replies[1]); 6104 6105 assert_eq!(accepted_count(&replies), if join_accepted { 2 } else { 1 }); 6106 if join_accepted { 6107 assert!(rejected_messages(&replies).is_empty()); 6108 } else { 6109 assert_eq!( 6110 rejected_messages(&replies), 6111 vec!["duplicate: group member already exists".to_owned()] 6112 ); 6113 } 6114 assert_runtime_member_status( 6115 &handle, 6116 "RaceJoinLeave", 6117 &FixtureKey::Member.public_key(), 6118 if join_accepted { 6119 MemberStatus::Member 6120 } else { 6121 MemberStatus::Removed 6122 }, 6123 ); 6124 assert_live_projection_matches_rebuild(&handle, "RaceJoinLeave"); 6125 6126 let _ = std::fs::remove_dir_all(root); 6127 } 6128 6129 #[tokio::test] 6130 async fn runtime_group_concurrency_delete_tombstone_blocks_normal_write() { 6131 let root = temp_root("runtime-group-concurrency-delete-write"); 6132 let _ = std::fs::remove_dir_all(&root); 6133 let handle = RelayRuntimeHandle::new( 6134 RelayRuntime::open(runtime_config(&root, 32)).expect("runtime"), 6135 ); 6136 let mut owner_auth = 6137 authenticated_runtime_state(&handle, FixtureKey::Owner, "owner-delete", 1_714_126_400) 6138 .await; 6139 let create = 6140 tangle_v2_group_create_event(FixtureKey::Owner, "RaceDelete", 1_714_126_401, &[]) 6141 .expect("create"); 6142 assert_accepted_reply( 6143 runtime_event_reply(&handle, create.clone(), &mut owner_auth, 1_714_126_401).await, 6144 &create, 6145 ); 6146 let normal = 6147 tangle_v2_group_event(FixtureKey::Owner, "RaceDelete", 1_714_126_402, 1, "normal") 6148 .expect("normal"); 6149 let delete = tangle_v2_delete_group_event(FixtureKey::Owner, "RaceDelete", 1_714_126_403) 6150 .expect("delete"); 6151 let normal_task = { 6152 let handle = handle.clone(); 6153 let mut auth = owner_auth.clone(); 6154 let event = normal.clone(); 6155 tokio::spawn(async move { 6156 runtime_event_reply(&handle, event, &mut auth, 1_714_126_402).await 6157 }) 6158 }; 6159 let delete_task = { 6160 let handle = handle.clone(); 6161 let mut auth = owner_auth.clone(); 6162 let event = delete.clone(); 6163 tokio::spawn(async move { 6164 runtime_event_reply(&handle, event, &mut auth, 1_714_126_403).await 6165 }) 6166 }; 6167 let replies = tokio::time::timeout(Duration::from_secs(3), async { 6168 vec![ 6169 normal_task.await.expect("normal task"), 6170 delete_task.await.expect("delete task"), 6171 ] 6172 }) 6173 .await 6174 .expect("delete write race"); 6175 let delete_reply = &replies[1]; 6176 6177 assert!(reply_is_accepted(delete_reply)); 6178 assert!( 6179 reply_is_accepted(&replies[0]) 6180 || rejected_messages(&replies) == vec!["blocked: group is deleted".to_owned()] 6181 ); 6182 let mut auth = owner_auth.clone(); 6183 assert_eq!( 6184 runtime_group_count( 6185 &handle, 6186 "deleted-normal-count", 6187 "RaceDelete", 6188 1, 6189 "h", 6190 &mut auth, 6191 1_714_126_404, 6192 ) 6193 .await, 6194 0 6195 ); 6196 assert_eq!( 6197 runtime_group_count( 6198 &handle, 6199 "deleted-marker-count", 6200 "RaceDelete", 6201 KIND_GROUP_DELETE_GROUP, 6202 "h", 6203 &mut auth, 6204 1_714_126_405, 6205 ) 6206 .await, 6207 1 6208 ); 6209 assert_live_projection_matches_rebuild(&handle, "RaceDelete"); 6210 6211 let _ = std::fs::remove_dir_all(root); 6212 } 6213 6214 #[tokio::test] 6215 async fn runtime_group_concurrency_membership_mutation_matches_rebuild() { 6216 let root = temp_root("runtime-group-concurrency-membership-mutation"); 6217 let _ = std::fs::remove_dir_all(&root); 6218 let handle = RelayRuntimeHandle::new( 6219 RelayRuntime::open(runtime_config(&root, 32)).expect("runtime"), 6220 ); 6221 let mut owner_auth = authenticated_runtime_state( 6222 &handle, 6223 FixtureKey::Owner, 6224 "owner-membership", 6225 1_714_126_500, 6226 ) 6227 .await; 6228 let create = 6229 tangle_v2_group_create_event(FixtureKey::Owner, "RaceMember", 1_714_126_501, &[]) 6230 .expect("create"); 6231 assert_accepted_reply( 6232 runtime_event_reply(&handle, create.clone(), &mut owner_auth, 1_714_126_501).await, 6233 &create, 6234 ); 6235 let put_member = tangle_v2_put_user_event( 6236 FixtureKey::Owner, 6237 "RaceMember", 6238 FixtureKey::Member, 6239 1_714_126_502, 6240 ) 6241 .expect("put member"); 6242 let remove_member = tangle_v2_remove_user_event( 6243 FixtureKey::Owner, 6244 "RaceMember", 6245 FixtureKey::Member, 6246 1_714_126_503, 6247 ) 6248 .expect("remove member"); 6249 let put_task = { 6250 let handle = handle.clone(); 6251 let mut auth = owner_auth.clone(); 6252 let event = put_member.clone(); 6253 tokio::spawn(async move { 6254 runtime_event_reply(&handle, event, &mut auth, 1_714_126_502).await 6255 }) 6256 }; 6257 let remove_task = { 6258 let handle = handle.clone(); 6259 let mut auth = owner_auth.clone(); 6260 let event = remove_member.clone(); 6261 tokio::spawn(async move { 6262 runtime_event_reply(&handle, event, &mut auth, 1_714_126_503).await 6263 }) 6264 }; 6265 let replies = tokio::time::timeout(Duration::from_secs(3), async { 6266 vec![ 6267 put_task.await.expect("put task"), 6268 remove_task.await.expect("remove task"), 6269 ] 6270 }) 6271 .await 6272 .expect("membership mutation race"); 6273 let remove_accepted = reply_is_accepted(&replies[1]); 6274 6275 assert!(reply_is_accepted(&replies[0])); 6276 if remove_accepted { 6277 assert!(rejected_messages(&replies).is_empty()); 6278 } else { 6279 assert_eq!( 6280 rejected_messages(&replies), 6281 vec!["duplicate: group member does not exist".to_owned()] 6282 ); 6283 } 6284 assert_runtime_member_status( 6285 &handle, 6286 "RaceMember", 6287 &FixtureKey::Member.public_key(), 6288 if remove_accepted { 6289 MemberStatus::Removed 6290 } else { 6291 MemberStatus::Member 6292 }, 6293 ); 6294 assert_live_projection_matches_rebuild(&handle, "RaceMember"); 6295 6296 let _ = std::fs::remove_dir_all(root); 6297 } 6298 6299 #[tokio::test] 6300 async fn runtime_shared_services_progress_under_concurrent_event_query_count_and_fanout() { 6301 let root = temp_root("runtime-shared-concurrency"); 6302 let _ = std::fs::remove_dir_all(&root); 6303 let handle = RelayRuntimeHandle::new( 6304 RelayRuntime::open(runtime_config(&root, 32)).expect("runtime"), 6305 ); 6306 let base_time = 1_714_126_000; 6307 let mut owner_auth = handle.auth_state().await.expect("owner auth"); 6308 owner_auth 6309 .issue_challenge("owner-stress", UnixTimestamp::new(base_time)) 6310 .expect("owner challenge"); 6311 let owner_auth_event = 6312 runtime_pocket_auth_event(FixtureKey::Owner, "owner-stress", base_time); 6313 assert_eq!( 6314 handle 6315 .handle_client_message( 6316 RuntimeClientMessage::Auth(owner_auth_event.clone()), 6317 &mut owner_auth, 6318 UnixTimestamp::new(base_time) 6319 ) 6320 .await 6321 .expect("owner auth"), 6322 vec![RelayMessage::Ok { 6323 event_id: runtime_pocket_event_id(&owner_auth_event), 6324 accepted: true, 6325 message: String::new() 6326 }] 6327 ); 6328 let create = runtime_pocket_group_create_event( 6329 FixtureKey::Owner, 6330 "StressPrivate", 6331 base_time + 1, 6332 &["private"], 6333 ); 6334 assert_eq!( 6335 handle 6336 .handle_client_message( 6337 RuntimeClientMessage::Event(create.clone()), 6338 &mut owner_auth, 6339 UnixTimestamp::new(base_time + 1) 6340 ) 6341 .await 6342 .expect("create"), 6343 vec![RelayMessage::Ok { 6344 event_id: runtime_pocket_event_id(&create), 6345 accepted: true, 6346 message: String::new() 6347 }] 6348 ); 6349 let put_member = runtime_pocket_put_user_event( 6350 FixtureKey::Owner, 6351 "StressPrivate", 6352 FixtureKey::Member, 6353 base_time + 2, 6354 ); 6355 assert_eq!( 6356 handle 6357 .handle_client_message( 6358 RuntimeClientMessage::Event(put_member.clone()), 6359 &mut owner_auth, 6360 UnixTimestamp::new(base_time + 2) 6361 ) 6362 .await 6363 .expect("put member"), 6364 vec![RelayMessage::Ok { 6365 event_id: runtime_pocket_event_id(&put_member), 6366 accepted: true, 6367 message: String::new() 6368 }] 6369 ); 6370 let mut member_auth = handle.auth_state().await.expect("member auth"); 6371 member_auth 6372 .issue_challenge("member-stress", UnixTimestamp::new(base_time + 3)) 6373 .expect("member challenge"); 6374 let member_auth_event = 6375 runtime_pocket_auth_event(FixtureKey::Member, "member-stress", base_time + 3); 6376 assert_eq!( 6377 handle 6378 .handle_client_message( 6379 RuntimeClientMessage::Auth(member_auth_event.clone()), 6380 &mut member_auth, 6381 UnixTimestamp::new(base_time + 3) 6382 ) 6383 .await 6384 .expect("member auth"), 6385 vec![RelayMessage::Ok { 6386 event_id: runtime_pocket_event_id(&member_auth_event), 6387 accepted: true, 6388 message: String::new() 6389 }] 6390 ); 6391 let public_auth = handle.auth_state().await.expect("public auth"); 6392 let mut offsets = handle.subscribe_events().await; 6393 let group_write_count = 6_usize; 6394 let public_write_count = 4_usize; 6395 let mut write_tasks = Vec::new(); 6396 for index in 0..group_write_count { 6397 let handle = handle.clone(); 6398 let mut auth = member_auth.clone(); 6399 write_tasks.push(tokio::spawn(async move { 6400 let event = runtime_pocket_group_event( 6401 FixtureKey::Member, 6402 "StressPrivate", 6403 base_time + 10 + u64::try_from(index).expect("index"), 6404 1, 6405 &format!("private stress {index}"), 6406 ); 6407 assert_eq!( 6408 handle 6409 .handle_client_message( 6410 RuntimeClientMessage::Event(event.clone()), 6411 &mut auth, 6412 UnixTimestamp::new( 6413 base_time + 10 + u64::try_from(index).expect("index") 6414 ) 6415 ) 6416 .await 6417 .expect("group write"), 6418 vec![RelayMessage::Ok { 6419 event_id: runtime_pocket_event_id(&event), 6420 accepted: true, 6421 message: String::new() 6422 }] 6423 ); 6424 (true, runtime_pocket_event_id(&event)) 6425 })); 6426 } 6427 for index in 0..public_write_count { 6428 let handle = handle.clone(); 6429 let mut auth = public_auth.clone(); 6430 write_tasks.push(tokio::spawn(async move { 6431 let event = runtime_pocket_event( 6432 FixtureKey::Admin, 6433 base_time + 40 + u64::try_from(index).expect("index"), 6434 1, 6435 Vec::new(), 6436 &format!("public stress {index}"), 6437 ); 6438 assert_eq!( 6439 handle 6440 .handle_client_message( 6441 RuntimeClientMessage::Event(event.clone()), 6442 &mut auth, 6443 UnixTimestamp::new( 6444 base_time + 40 + u64::try_from(index).expect("index") 6445 ) 6446 ) 6447 .await 6448 .expect("public write"), 6449 vec![RelayMessage::Ok { 6450 event_id: runtime_pocket_event_id(&event), 6451 accepted: true, 6452 message: String::new() 6453 }] 6454 ); 6455 (false, runtime_pocket_event_id(&event)) 6456 })); 6457 } 6458 let stored_events = tokio::time::timeout(Duration::from_secs(3), async { 6459 let mut stored_events = Vec::new(); 6460 for task in write_tasks { 6461 stored_events.push(task.await.expect("write task")); 6462 } 6463 stored_events 6464 }) 6465 .await 6466 .expect("write concurrency timeout"); 6467 assert_eq!( 6468 stored_events 6469 .iter() 6470 .filter(|(is_group, _)| *is_group) 6471 .count(), 6472 group_write_count 6473 ); 6474 assert_eq!( 6475 stored_events 6476 .iter() 6477 .filter(|(is_group, _)| !*is_group) 6478 .count(), 6479 public_write_count 6480 ); 6481 let group_event_ids = stored_events 6482 .iter() 6483 .filter(|(is_group, _)| *is_group) 6484 .map(|(_, event_id)| event_id.clone()) 6485 .collect::<BTreeSet<_>>(); 6486 let mut published_offsets = Vec::new(); 6487 for _ in 0..stored_events.len() { 6488 published_offsets.push( 6489 tokio::time::timeout(Duration::from_secs(1), offsets.recv()) 6490 .await 6491 .expect("offset timeout") 6492 .expect("offset"), 6493 ); 6494 } 6495 assert_eq!( 6496 offsets.try_recv().expect_err("no extra offsets"), 6497 TangleEventReceiveError::Empty 6498 ); 6499 let mut visibility_tasks = Vec::new(); 6500 for offset in published_offsets.iter().copied() { 6501 let handle = handle.clone(); 6502 let member_auth = member_auth.clone(); 6503 let public_auth = public_auth.clone(); 6504 let group_event_ids = group_event_ids.clone(); 6505 visibility_tasks.push(tokio::spawn(async move { 6506 let member_event = handle 6507 .event_by_offset_with_auth(offset, &member_auth) 6508 .await 6509 .expect("member offset") 6510 .expect("member visible"); 6511 let public_event = handle 6512 .event_by_offset_with_auth(offset, &public_auth) 6513 .await 6514 .expect("public offset"); 6515 let member_event_id = 6516 EventId::new(&member_event.id().as_hex_string()).expect("pocket id"); 6517 let is_group_event = group_event_ids.contains(&member_event_id); 6518 if is_group_event { 6519 assert!(public_event.is_none()); 6520 } else { 6521 assert!(public_event.is_some()); 6522 } 6523 is_group_event 6524 })); 6525 } 6526 let visible_group_offsets = tokio::time::timeout(Duration::from_secs(3), async { 6527 let mut visible_group_offsets = 0; 6528 for task in visibility_tasks { 6529 if task.await.expect("visibility task") { 6530 visible_group_offsets += 1; 6531 } 6532 } 6533 visible_group_offsets 6534 }) 6535 .await 6536 .expect("visibility timeout"); 6537 assert_eq!(visible_group_offsets, group_write_count); 6538 let member_subscription = SubscriptionId::new("member-stress-live").expect("subscription"); 6539 let public_subscription = SubscriptionId::new("public-stress-live").expect("subscription"); 6540 let mut member_subscriptions = LiveSubscriptionSet::new(32, 64).expect("member live set"); 6541 let mut public_subscriptions = LiveSubscriptionSet::new(32, 64).expect("public live set"); 6542 let stress_filter = pocket_filter(json!({"kinds":[1], "#h":["StressPrivate"]})); 6543 member_subscriptions 6544 .subscribe(member_subscription.clone(), vec![stress_filter.clone()]) 6545 .expect("member subscribe"); 6546 public_subscriptions 6547 .subscribe(public_subscription, vec![stress_filter]) 6548 .expect("public subscribe"); 6549 let mut member_fanout_count = 0; 6550 for offset in &published_offsets { 6551 let public_replies = handle 6552 .fanout_event_offset_with_projection_context( 6553 *offset, 6554 &mut public_subscriptions, 6555 &public_auth, 6556 &RelayProjectionContext::default(), 6557 ) 6558 .await 6559 .expect("public fanout"); 6560 assert!(public_replies.is_empty()); 6561 let member_replies = handle 6562 .fanout_event_offset_with_projection_context( 6563 *offset, 6564 &mut member_subscriptions, 6565 &member_auth, 6566 &RelayProjectionContext::default(), 6567 ) 6568 .await 6569 .expect("member fanout"); 6570 for reply in member_replies { 6571 match reply { 6572 RuntimeRelayMessage::Event { 6573 subscription_id, 6574 event, 6575 } => { 6576 assert_eq!(subscription_id, member_subscription); 6577 let event_id = 6578 EventId::new(&event.id().as_hex_string()).expect("pocket id"); 6579 assert!(group_event_ids.contains(&event_id)); 6580 member_fanout_count += 1; 6581 } 6582 other => panic!("unexpected fanout reply {other:?}"), 6583 } 6584 } 6585 } 6586 assert_eq!(member_fanout_count, group_write_count); 6587 let mut query_tasks = Vec::new(); 6588 for index in 0..3_u64 { 6589 let member_req_handle = handle.clone(); 6590 let mut auth = member_auth.clone(); 6591 let group_event_ids = group_event_ids.clone(); 6592 query_tasks.push(tokio::spawn(async move { 6593 let subscription_id = 6594 SubscriptionId::new(&format!("member-req-{index}")).expect("subscription"); 6595 let replies = member_req_handle 6596 .handle_protocol_client_message_for_test( 6597 ClientMessage::Req { 6598 subscription_id: subscription_id.clone(), 6599 filters: vec![ 6600 filter_from_value(&json!({ 6601 "kinds":[1], 6602 "#h":["StressPrivate"], 6603 "limit": 20 6604 })) 6605 .expect("filter"), 6606 ], 6607 }, 6608 &mut auth, 6609 UnixTimestamp::new(base_time + 100 + index), 6610 ) 6611 .await 6612 .expect("member req"); 6613 assert_eq!( 6614 replies 6615 .iter() 6616 .filter(|reply| matches!( 6617 reply, 6618 RelayMessage::Event { 6619 subscription_id: delivered, 6620 event 6621 } if delivered == &subscription_id && group_event_ids.contains(event.id()) 6622 )) 6623 .count(), 6624 group_event_ids.len() 6625 ); 6626 assert!(matches!( 6627 replies.last(), 6628 Some(RelayMessage::Eose(delivered)) if delivered == &subscription_id 6629 )); 6630 })); 6631 let public_req_handle = handle.clone(); 6632 let mut auth = public_auth.clone(); 6633 query_tasks.push(tokio::spawn(async move { 6634 let subscription_id = 6635 SubscriptionId::new(&format!("public-req-{index}")).expect("subscription"); 6636 let replies = public_req_handle 6637 .handle_protocol_client_message_for_test( 6638 ClientMessage::Req { 6639 subscription_id: subscription_id.clone(), 6640 filters: vec![ 6641 filter_from_value(&json!({ 6642 "kinds":[1], 6643 "#h":["StressPrivate"], 6644 "limit": 20 6645 })) 6646 .expect("filter"), 6647 ], 6648 }, 6649 &mut auth, 6650 UnixTimestamp::new(base_time + 110 + index), 6651 ) 6652 .await 6653 .expect("public req"); 6654 assert_eq!( 6655 replies, 6656 vec![RelayMessage::Closed { 6657 subscription_id, 6658 message: "auth-required: authentication required to read group events" 6659 .to_owned() 6660 }] 6661 ); 6662 })); 6663 let member_count_handle = handle.clone(); 6664 let mut auth = member_auth.clone(); 6665 query_tasks.push(tokio::spawn(async move { 6666 let subscription_id = 6667 SubscriptionId::new(&format!("member-count-{index}")).expect("subscription"); 6668 let replies = member_count_handle 6669 .handle_protocol_client_message_for_test( 6670 ClientMessage::Count { 6671 subscription_id: subscription_id.clone(), 6672 filters: vec![ 6673 filter_from_value(&json!({ 6674 "kinds":[1], 6675 "#h":["StressPrivate"] 6676 })) 6677 .expect("filter"), 6678 ], 6679 }, 6680 &mut auth, 6681 UnixTimestamp::new(base_time + 120 + index), 6682 ) 6683 .await 6684 .expect("member count"); 6685 assert_eq!( 6686 replies, 6687 vec![RelayMessage::Count { 6688 subscription_id, 6689 count: u64::try_from(group_write_count).expect("group count"), 6690 hll: None 6691 }] 6692 ); 6693 })); 6694 let public_count_handle = handle.clone(); 6695 let mut auth = public_auth.clone(); 6696 query_tasks.push(tokio::spawn(async move { 6697 let subscription_id = 6698 SubscriptionId::new(&format!("public-count-{index}")).expect("subscription"); 6699 let replies = public_count_handle 6700 .handle_protocol_client_message_for_test( 6701 ClientMessage::Count { 6702 subscription_id: subscription_id.clone(), 6703 filters: vec![ 6704 filter_from_value(&json!({ 6705 "kinds":[1], 6706 "#h":["StressPrivate"] 6707 })) 6708 .expect("filter"), 6709 ], 6710 }, 6711 &mut auth, 6712 UnixTimestamp::new(base_time + 130 + index), 6713 ) 6714 .await 6715 .expect("public count"); 6716 assert_eq!( 6717 replies, 6718 vec![RelayMessage::Count { 6719 subscription_id, 6720 count: 0, 6721 hll: None 6722 }] 6723 ); 6724 })); 6725 } 6726 tokio::time::timeout(Duration::from_secs(3), async { 6727 for task in query_tasks { 6728 task.await.expect("query task"); 6729 } 6730 }) 6731 .await 6732 .expect("query concurrency timeout"); 6733 assert!(handle.metrics().query_candidates_scanned() > 0); 6734 assert!( 6735 handle.metrics().query_returned_events() 6736 >= u64::try_from(group_write_count * 3).expect("returned event count") 6737 ); 6738 assert!(handle.metrics().query_redacted_events() > 0); 6739 handle.shutdown().await.expect("shutdown"); 6740 6741 let _ = std::fs::remove_dir_all(root); 6742 } 6743 6744 fn runtime_config(root: &Path, per_connection_outbound_queue: usize) -> BaseRelayRuntimeConfig { 6745 runtime_config_with_group_policy(root, per_connection_outbound_queue, false) 6746 } 6747 6748 fn runtime_config_with_public_join( 6749 root: &Path, 6750 per_connection_outbound_queue: usize, 6751 ) -> BaseRelayRuntimeConfig { 6752 runtime_config_with_group_policy(root, per_connection_outbound_queue, true) 6753 } 6754 6755 #[derive(Default)] 6756 struct RecordingHooks { 6757 admissions: Mutex<Vec<RelayEventAdmissionContext>>, 6758 stored: Mutex<Vec<RelayEventStoredContext>>, 6759 stored_bytes: std::sync::atomic::AtomicU64, 6760 } 6761 6762 impl RelayRuntimeHooks for RecordingHooks { 6763 fn admit_event(&self, context: &RelayEventAdmissionContext) -> EventAdmissionDecision { 6764 self.admissions 6765 .lock() 6766 .expect("admissions") 6767 .push(context.clone()); 6768 if context.event().has_tag("policy", "reject") { 6769 EventAdmissionDecision::reject("hook rejected event") 6770 } else { 6771 EventAdmissionDecision::Accept 6772 } 6773 } 6774 6775 fn event_stored(&self, context: &RelayEventStoredContext) { 6776 self.stored.lock().expect("stored").push(context.clone()); 6777 self.stored_bytes 6778 .store(144, std::sync::atomic::Ordering::Relaxed); 6779 } 6780 6781 fn storage_used_bytes(&self) -> Option<u64> { 6782 Some(self.stored_bytes.load(std::sync::atomic::Ordering::Relaxed)) 6783 } 6784 } 6785 6786 struct BlockingStoredHooks { 6787 started_sender: Mutex<Option<mpsc::SyncSender<()>>>, 6788 release_receiver: Mutex<mpsc::Receiver<()>>, 6789 } 6790 6791 impl RelayRuntimeHooks for BlockingStoredHooks { 6792 fn event_stored(&self, _context: &RelayEventStoredContext) { 6793 if let Some(sender) = self.started_sender.lock().expect("started sender").take() { 6794 sender.send(()).expect("send hook start"); 6795 } 6796 self.release_receiver 6797 .lock() 6798 .expect("release receiver") 6799 .recv() 6800 .expect("receive hook release"); 6801 } 6802 } 6803 6804 #[derive(Debug, Clone, Copy, PartialEq, Eq)] 6805 enum ProjectionHookScope { 6806 Historical, 6807 Live, 6808 } 6809 6810 struct ProjectingHooks { 6811 projection_identifier: &'static str, 6812 scope: ProjectionHookScope, 6813 source_content: Option<&'static str>, 6814 decision: Mutex<RelayEventProjectionDecision>, 6815 query_plan: Mutex<RelayProjectionQueryPlan>, 6816 live_candidates: Mutex<Vec<RelayLiveProjectionCandidate>>, 6817 contexts: Mutex<Vec<RelayEventProjectionContext>>, 6818 live_contexts: Mutex<Vec<RelayLiveProjectionContext>>, 6819 query_contexts: Mutex<Vec<RelayQueryProjectionContext>>, 6820 } 6821 6822 impl ProjectingHooks { 6823 fn new( 6824 projection_identifier: &'static str, 6825 scope: ProjectionHookScope, 6826 source_content: Option<&'static str>, 6827 decision: RelayEventProjectionDecision, 6828 ) -> Self { 6829 Self { 6830 projection_identifier, 6831 scope, 6832 source_content, 6833 decision: Mutex::new(decision), 6834 query_plan: Mutex::new(RelayProjectionQueryPlan::default()), 6835 live_candidates: Mutex::new(Vec::new()), 6836 contexts: Mutex::new(Vec::new()), 6837 live_contexts: Mutex::new(Vec::new()), 6838 query_contexts: Mutex::new(Vec::new()), 6839 } 6840 } 6841 6842 fn set_decision(&self, decision: RelayEventProjectionDecision) { 6843 *self.decision.lock().expect("decision") = decision; 6844 } 6845 6846 fn set_query_plan(&self, plan: RelayProjectionQueryPlan) { 6847 *self.query_plan.lock().expect("query plan") = plan; 6848 } 6849 6850 fn set_live_candidates(&self, candidates: Vec<RelayLiveProjectionCandidate>) { 6851 *self.live_candidates.lock().expect("live candidates") = candidates; 6852 } 6853 6854 fn contexts(&self) -> Vec<RelayEventProjectionContext> { 6855 self.contexts.lock().expect("contexts").clone() 6856 } 6857 6858 fn live_contexts(&self) -> Vec<RelayLiveProjectionContext> { 6859 self.live_contexts.lock().expect("live contexts").clone() 6860 } 6861 6862 fn query_contexts(&self) -> Vec<RelayQueryProjectionContext> { 6863 self.query_contexts.lock().expect("query contexts").clone() 6864 } 6865 6866 fn scope_matches(&self, source: RelayEventProjectionSource) -> bool { 6867 matches!( 6868 (self.scope, source), 6869 ( 6870 ProjectionHookScope::Historical, 6871 RelayEventProjectionSource::HistoricalQuery 6872 ) | ( 6873 ProjectionHookScope::Live, 6874 RelayEventProjectionSource::LiveFanout { .. } 6875 ) 6876 ) 6877 } 6878 6879 fn content_matches(&self, context: &RelayEventProjectionContext) -> bool { 6880 match self.source_content { 6881 Some(content) => context.event().content() == content, 6882 None => true, 6883 } 6884 } 6885 } 6886 6887 impl RelayRuntimeHooks for ProjectingHooks { 6888 fn plan_query(&self, context: &RelayQueryProjectionContext) -> RelayProjectionQueryPlan { 6889 self.query_contexts 6890 .lock() 6891 .expect("query contexts") 6892 .push(context.clone()); 6893 *self.query_plan.lock().expect("query plan") 6894 } 6895 6896 fn live_projection_candidates( 6897 &self, 6898 context: &RelayLiveProjectionContext, 6899 ) -> Vec<RelayLiveProjectionCandidate> { 6900 self.live_contexts 6901 .lock() 6902 .expect("live contexts") 6903 .push(context.clone()); 6904 if context.projection().identifier() == Some(self.projection_identifier) { 6905 self.live_candidates 6906 .lock() 6907 .expect("live candidates") 6908 .clone() 6909 } else { 6910 Vec::new() 6911 } 6912 } 6913 6914 fn project_event( 6915 &self, 6916 context: &RelayEventProjectionContext, 6917 ) -> RelayEventProjectionDecision { 6918 self.contexts 6919 .lock() 6920 .expect("contexts") 6921 .push(context.clone()); 6922 if context.projection().identifier() == Some(self.projection_identifier) 6923 && self.scope_matches(context.source()) 6924 && self.content_matches(context) 6925 { 6926 *self.decision.lock().expect("decision") 6927 } else { 6928 RelayEventProjectionDecision::Emit 6929 } 6930 } 6931 } 6932 6933 fn runtime_config_with_group_policy( 6934 root: &Path, 6935 per_connection_outbound_queue: usize, 6936 public_join: bool, 6937 ) -> BaseRelayRuntimeConfig { 6938 let raw = json!({ 6939 "server": { 6940 "listen_addr": "127.0.0.1:0", 6941 "relay_url": "wss://relay.radroots.test" 6942 }, 6943 "pocket": { 6944 "data_directory": root.join("pocket"), 6945 "sync_policy": "flush_on_shutdown", 6946 "query": { 6947 "allow_scraping": false, 6948 "allow_scrape_if_limited_to": 100, 6949 "allow_scrape_if_max_seconds": 3600 6950 } 6951 }, 6952 "groups": { 6953 "enabled": true, 6954 "canonical_relay_url": "wss://relay.radroots.test", 6955 "relay_secret": "7777777777777777777777777777777777777777777777777777777777777777", 6956 "owner_pubkeys": [FixtureKey::Owner.public_key().as_str()], 6957 "policy": { 6958 "public_join": public_join, 6959 "invites_enabled": false 6960 } 6961 }, 6962 "auth": { 6963 "challenge_ttl_seconds": 300, 6964 "created_at_skew_seconds": 600 6965 }, 6966 "limits": { 6967 "max_message_length": 1048576, 6968 "max_subid_length": 64, 6969 "max_subscriptions_per_connection": 64, 6970 "max_filters_per_request": 10, 6971 "max_tag_values_per_filter": 100, 6972 "max_query_complexity": 2048, 6973 "max_limit": 500, 6974 "default_limit": 100, 6975 "max_event_tags": 200, 6976 "max_content_length": 65536, 6977 "broadcast_channel_capacity": 16, 6978 "per_connection_outbound_queue": per_connection_outbound_queue 6979 }, 6980 "rate_limits": { 6981 "auth": { 6982 "per_ip": {"window_seconds": 60, "max_hits": 120}, 6983 "per_pubkey": {"window_seconds": 60, "max_hits": 30}, 6984 "failures": {"window_seconds": 300, "max_hits": 5}, 6985 "failures_per_ip": {"window_seconds": 300, "max_hits": 20} 6986 }, 6987 "event": { 6988 "per_ip": {"window_seconds": 60, "max_hits": 600}, 6989 "per_pubkey": {"window_seconds": 60, "max_hits": 120}, 6990 "per_kind": {"window_seconds": 60, "max_hits": 1000} 6991 }, 6992 "group": { 6993 "write_per_ip": {"window_seconds": 60, "max_hits": 300}, 6994 "write_per_pubkey": {"window_seconds": 60, "max_hits": 60}, 6995 "write_per_group": {"window_seconds": 60, "max_hits": 90}, 6996 "write_per_kind": {"window_seconds": 60, "max_hits": 300}, 6997 "join_flow": {"window_seconds": 300, "max_hits": 10}, 6998 "join_flow_per_ip": {"window_seconds": 300, "max_hits": 30} 6999 }, 7000 "req": { 7001 "per_ip": {"window_seconds": 60, "max_hits": 600}, 7002 "per_connection": {"window_seconds": 60, "max_hits": 120}, 7003 "per_pubkey": {"window_seconds": 60, "max_hits": 240}, 7004 "per_group": {"window_seconds": 60, "max_hits": 240}, 7005 "per_kind": {"window_seconds": 60, "max_hits": 500}, 7006 "broad": {"window_seconds": 60, "max_hits": 30} 7007 }, 7008 "count": { 7009 "per_ip": {"window_seconds": 60, "max_hits": 300}, 7010 "per_connection": {"window_seconds": 60, "max_hits": 60}, 7011 "per_pubkey": {"window_seconds": 60, "max_hits": 120}, 7012 "per_group": {"window_seconds": 60, "max_hits": 120}, 7013 "per_kind": {"window_seconds": 60, "max_hits": 240}, 7014 "broad": {"window_seconds": 60, "max_hits": 20} 7015 } 7016 } 7017 }) 7018 .to_string(); 7019 parse_base_relay_runtime_config_json(&raw).expect("config") 7020 } 7021 7022 async fn authenticated_runtime_state( 7023 handle: &RelayRuntimeHandle, 7024 key: FixtureKey, 7025 challenge: &str, 7026 now: u64, 7027 ) -> BaseAuthState { 7028 let mut auth = handle.auth_state().await.expect("auth"); 7029 auth.issue_challenge(challenge, UnixTimestamp::new(now)) 7030 .expect("challenge"); 7031 let event = tangle_v2_auth_event(key, challenge, now).expect("auth event"); 7032 let replies = handle 7033 .handle_protocol_client_message_for_test( 7034 ClientMessage::Auth(event.clone()), 7035 &mut auth, 7036 UnixTimestamp::new(now), 7037 ) 7038 .await 7039 .expect("auth message"); 7040 7041 assert_eq!( 7042 replies, 7043 vec![RelayMessage::Ok { 7044 event_id: event.id().clone(), 7045 accepted: true, 7046 message: String::new() 7047 }] 7048 ); 7049 auth 7050 } 7051 7052 async fn runtime_event_reply( 7053 handle: &RelayRuntimeHandle, 7054 event: Event, 7055 auth: &mut BaseAuthState, 7056 now: u64, 7057 ) -> RelayMessage { 7058 let replies = handle 7059 .handle_protocol_client_message_for_test( 7060 ClientMessage::Event(event), 7061 auth, 7062 UnixTimestamp::new(now), 7063 ) 7064 .await 7065 .expect("event message"); 7066 7067 assert_eq!(replies.len(), 1); 7068 replies.into_iter().next().expect("reply") 7069 } 7070 7071 fn runtime_pocket_event_reply( 7072 handle: &RelayRuntimeHandle, 7073 event: &PocketEvent, 7074 auth: &mut BaseAuthState, 7075 ) -> RelayMessage { 7076 handle 7077 .inner 7078 .handle_pocket_event_with_auth_report(event, auth) 7079 .expect("event message") 7080 .into_message() 7081 } 7082 7083 async fn runtime_group_count( 7084 handle: &RelayRuntimeHandle, 7085 subscription_id: &str, 7086 group_id: &str, 7087 kind: u32, 7088 tag_name: &str, 7089 auth: &mut BaseAuthState, 7090 now: u64, 7091 ) -> u64 { 7092 let replies = handle 7093 .handle_protocol_client_message_for_test( 7094 ClientMessage::Count { 7095 subscription_id: SubscriptionId::new(subscription_id).expect("subscription"), 7096 filters: vec![runtime_group_filter(group_id, kind, tag_name)], 7097 }, 7098 auth, 7099 UnixTimestamp::new(now), 7100 ) 7101 .await 7102 .expect("count"); 7103 7104 match replies.as_slice() { 7105 [RelayMessage::Count { count, .. }] => *count, 7106 other => panic!("count reply expected, got {other:?}"), 7107 } 7108 } 7109 7110 fn runtime_group_filter(group_id: &str, kind: u32, tag_name: &str) -> Filter { 7111 let mut value = json!({"kinds": [kind]}); 7112 value 7113 .as_object_mut() 7114 .expect("filter") 7115 .insert(format!("#{tag_name}"), json!([group_id])); 7116 filter_from_value(&value).expect("filter") 7117 } 7118 7119 async fn drain_offsets(receiver: &mut TangleEventReceiver, count: usize) -> Vec<StoreOffset> { 7120 let mut offsets = Vec::with_capacity(count); 7121 for _ in 0..count { 7122 offsets.push( 7123 tokio::time::timeout(Duration::from_secs(1), receiver.recv()) 7124 .await 7125 .expect("offset timeout") 7126 .expect("offset"), 7127 ); 7128 } 7129 offsets 7130 } 7131 7132 fn accepted_count(replies: &[RelayMessage]) -> usize { 7133 replies 7134 .iter() 7135 .filter(|reply| reply_is_accepted(reply)) 7136 .count() 7137 } 7138 7139 fn reply_is_accepted(reply: &RelayMessage) -> bool { 7140 matches!( 7141 reply, 7142 RelayMessage::Ok { 7143 accepted: true, 7144 message, 7145 .. 7146 } if message.is_empty() 7147 ) 7148 } 7149 7150 fn rejected_messages(replies: &[RelayMessage]) -> Vec<String> { 7151 replies 7152 .iter() 7153 .filter_map(|reply| match reply { 7154 RelayMessage::Ok { 7155 accepted: false, 7156 message, 7157 .. 7158 } => Some(message.clone()), 7159 _ => None, 7160 }) 7161 .collect() 7162 } 7163 7164 fn assert_accepted_reply(reply: RelayMessage, event: &Event) { 7165 assert_eq!( 7166 reply, 7167 RelayMessage::Ok { 7168 event_id: event.id().clone(), 7169 accepted: true, 7170 message: String::new() 7171 } 7172 ); 7173 } 7174 7175 fn assert_accepted_pocket_reply(reply: RelayMessage, event: &PocketEvent) { 7176 assert_eq!( 7177 reply, 7178 RelayMessage::Ok { 7179 event_id: runtime_pocket_event_id(event), 7180 accepted: true, 7181 message: String::new() 7182 } 7183 ); 7184 } 7185 7186 fn runtime_pocket_event_id(event: &PocketEvent) -> EventId { 7187 EventId::new(&event.id().as_hex_string()).expect("event id") 7188 } 7189 7190 fn assert_runtime_member_status( 7191 handle: &RelayRuntimeHandle, 7192 group_id: &str, 7193 pubkey: &PublicKeyHex, 7194 status: MemberStatus, 7195 ) { 7196 let group_id = GroupId::new(group_id).expect("group"); 7197 let groups = handle.inner.groups.as_ref().expect("groups"); 7198 let projection = groups.projection(); 7199 7200 assert_eq!( 7201 projection 7202 .member(&group_id, pubkey) 7203 .expect("member") 7204 .status(), 7205 status 7206 ); 7207 } 7208 7209 fn assert_live_projection_matches_rebuild(handle: &RelayRuntimeHandle, group_id: &str) { 7210 let group_id = GroupId::new(group_id).expect("group"); 7211 let groups = handle.inner.groups.as_ref().expect("groups"); 7212 let live = groups.projection(); 7213 let rebuilt = rebuilt_projection(handle); 7214 let live_group = live.group(&group_id); 7215 let rebuilt_group = rebuilt.group(&group_id); 7216 7217 assert_eq!( 7218 live_group.map(|group| group.lifecycle()), 7219 rebuilt_group.map(|group| group.lifecycle()) 7220 ); 7221 assert_eq!( 7222 live_group.map(|group| group.metadata()), 7223 rebuilt_group.map(|group| group.metadata()) 7224 ); 7225 assert_eq!( 7226 live_group.and_then(|group| group.delete_event_id()), 7227 rebuilt_group.and_then(|group| group.delete_event_id()) 7228 ); 7229 assert_eq!(live.tombstone(&group_id), rebuilt.tombstone(&group_id)); 7230 assert_eq!( 7231 projection_member_statuses(&live, &group_id), 7232 projection_member_statuses(&rebuilt, &group_id) 7233 ); 7234 } 7235 7236 fn rebuilt_projection(handle: &RelayRuntimeHandle) -> GroupProjection { 7237 let groups = handle.inner.groups.as_ref().expect("groups"); 7238 let limits = groups.limits(); 7239 let events = handle 7240 .inner 7241 .store 7242 .scan_events() 7243 .expect("scan") 7244 .into_iter() 7245 .filter_map(|stored| { 7246 let store_offset = StoreOffset::new(stored.store_offset()); 7247 match tangle_groups::classify_group_event(stored.event(), limits).expect("classify") 7248 { 7249 GroupEventClass::NonGroup => None, 7250 _ => Some(CanonicalGroupEvent::new(stored.into_event(), store_offset)), 7251 } 7252 }) 7253 .collect::<Vec<_>>(); 7254 7255 rebuild_group_projection(events, limits, UnixTimestamp::new(1_714_199_999)) 7256 .expect("rebuild") 7257 .into_projection() 7258 } 7259 7260 fn projection_member_statuses( 7261 projection: &GroupProjection, 7262 group_id: &GroupId, 7263 ) -> BTreeMap<String, MemberStatus> { 7264 projection 7265 .members() 7266 .iter() 7267 .filter(|((candidate, _), _)| candidate == group_id) 7268 .map(|((_, pubkey), member)| (pubkey.as_str().to_owned(), member.status())) 7269 .collect() 7270 } 7271 7272 fn runtime_relay_limits(max_pending_events: usize) -> BaseRelayLimits { 7273 BaseRelayLimits::new(BaseRelayLimitSettings { 7274 max_pending_events, 7275 max_subscription_id_length: 64, 7276 max_subscriptions: 64, 7277 max_filters_per_request: 10, 7278 max_tag_values_per_filter: 100, 7279 max_query_complexity: 610, 7280 max_event_tags: 200, 7281 max_content_length: 65_536, 7282 max_limit: 500, 7283 default_limit: 100, 7284 }) 7285 .expect("limits") 7286 } 7287 7288 fn pocket_filter(value: serde_json::Value) -> tangle_store_pocket::PocketOwnedFilter { 7289 let filter = filter_from_value(&value).expect("filter"); 7290 crate::pocket_conversion::tangle_filter_to_pocket(&filter).expect("pocket filter") 7291 } 7292 7293 fn tangle_v2_event( 7294 key: FixtureKey, 7295 created_at: u64, 7296 kind: u64, 7297 tags: Vec<Tag>, 7298 content: &str, 7299 ) -> Result<Event, String> { 7300 let event = runtime_pocket_event(key, created_at, kind, tags, content); 7301 runtime_pocket_event_to_protocol(&event) 7302 } 7303 7304 fn tangle_v2_auth_event( 7305 key: FixtureKey, 7306 challenge: &str, 7307 created_at: u64, 7308 ) -> Result<Event, String> { 7309 tangle_v2_event( 7310 key, 7311 created_at, 7312 22_242, 7313 vec![ 7314 Tag::from_parts("relay", &["wss://relay.radroots.test"])?, 7315 Tag::from_parts("challenge", &[challenge])?, 7316 ], 7317 "", 7318 ) 7319 } 7320 7321 fn tangle_v2_group_create_event( 7322 key: FixtureKey, 7323 group_id: &str, 7324 created_at: u64, 7325 flags: &[&str], 7326 ) -> Result<Event, String> { 7327 let mut tags = vec![ 7328 Tag::from_parts("h", &[group_id])?, 7329 Tag::from_parts("name", &[group_id])?, 7330 ]; 7331 for flag in flags { 7332 tags.push(Tag::from_parts(flag, &[])?); 7333 } 7334 tangle_v2_event(key, created_at, KIND_GROUP_CREATE_GROUP.into(), tags, "") 7335 } 7336 7337 fn tangle_v2_put_user_event( 7338 key: FixtureKey, 7339 group_id: &str, 7340 target: FixtureKey, 7341 created_at: u64, 7342 ) -> Result<Event, String> { 7343 let target_pubkey = target.public_key(); 7344 tangle_v2_event( 7345 key, 7346 created_at, 7347 KIND_GROUP_PUT_USER.into(), 7348 vec![ 7349 Tag::from_parts("h", &[group_id])?, 7350 Tag::from_parts("p", &[target_pubkey.as_str()])?, 7351 ], 7352 "", 7353 ) 7354 } 7355 7356 fn tangle_v2_remove_user_event( 7357 key: FixtureKey, 7358 group_id: &str, 7359 target: FixtureKey, 7360 created_at: u64, 7361 ) -> Result<Event, String> { 7362 let target_pubkey = target.public_key(); 7363 tangle_v2_event( 7364 key, 7365 created_at, 7366 KIND_GROUP_REMOVE_USER.into(), 7367 vec![ 7368 Tag::from_parts("h", &[group_id])?, 7369 Tag::from_parts("p", &[target_pubkey.as_str()])?, 7370 ], 7371 "", 7372 ) 7373 } 7374 7375 fn tangle_v2_join_event( 7376 key: FixtureKey, 7377 group_id: &str, 7378 created_at: u64, 7379 ) -> Result<Event, String> { 7380 tangle_v2_group_event( 7381 key, 7382 group_id, 7383 created_at, 7384 KIND_GROUP_JOIN_REQUEST.into(), 7385 "", 7386 ) 7387 } 7388 7389 fn tangle_v2_leave_event( 7390 key: FixtureKey, 7391 group_id: &str, 7392 created_at: u64, 7393 ) -> Result<Event, String> { 7394 tangle_v2_group_event( 7395 key, 7396 group_id, 7397 created_at, 7398 KIND_GROUP_LEAVE_REQUEST.into(), 7399 "", 7400 ) 7401 } 7402 7403 fn tangle_v2_delete_group_event( 7404 key: FixtureKey, 7405 group_id: &str, 7406 created_at: u64, 7407 ) -> Result<Event, String> { 7408 tangle_v2_group_event( 7409 key, 7410 group_id, 7411 created_at, 7412 KIND_GROUP_DELETE_GROUP.into(), 7413 "", 7414 ) 7415 } 7416 7417 fn tangle_v2_group_event( 7418 key: FixtureKey, 7419 group_id: &str, 7420 created_at: u64, 7421 kind: u64, 7422 content: &str, 7423 ) -> Result<Event, String> { 7424 tangle_v2_event( 7425 key, 7426 created_at, 7427 kind, 7428 vec![Tag::from_parts("h", &[group_id])?], 7429 content, 7430 ) 7431 } 7432 7433 fn runtime_pocket_group_create_event( 7434 key: FixtureKey, 7435 group_id: &str, 7436 created_at: u64, 7437 flags: &[&str], 7438 ) -> PocketOwnedEvent { 7439 let mut tags = vec![ 7440 Tag::from_parts("h", &[group_id]).expect("h"), 7441 Tag::from_parts("name", &[group_id]).expect("name"), 7442 ]; 7443 for flag in flags { 7444 tags.push(Tag::from_parts(flag, &[]).expect("flag")); 7445 } 7446 runtime_pocket_event(key, created_at, KIND_GROUP_CREATE_GROUP.into(), tags, "") 7447 } 7448 7449 fn runtime_pocket_auth_event( 7450 key: FixtureKey, 7451 challenge: &str, 7452 created_at: u64, 7453 ) -> PocketOwnedEvent { 7454 runtime_pocket_event( 7455 key, 7456 created_at, 7457 22_242, 7458 vec![ 7459 Tag::from_parts("relay", &["wss://relay.radroots.test"]).expect("relay"), 7460 Tag::from_parts("challenge", &[challenge]).expect("challenge"), 7461 ], 7462 "", 7463 ) 7464 } 7465 7466 fn runtime_pocket_put_user_event( 7467 key: FixtureKey, 7468 group_id: &str, 7469 target: FixtureKey, 7470 created_at: u64, 7471 ) -> PocketOwnedEvent { 7472 let target_pubkey = target.public_key(); 7473 runtime_pocket_event( 7474 key, 7475 created_at, 7476 KIND_GROUP_PUT_USER.into(), 7477 vec![ 7478 Tag::from_parts("h", &[group_id]).expect("h"), 7479 Tag::from_parts("p", &[target_pubkey.as_str()]).expect("p"), 7480 ], 7481 "", 7482 ) 7483 } 7484 7485 fn runtime_pocket_group_event( 7486 key: FixtureKey, 7487 group_id: &str, 7488 created_at: u64, 7489 kind: u64, 7490 content: &str, 7491 ) -> PocketOwnedEvent { 7492 runtime_pocket_event( 7493 key, 7494 created_at, 7495 kind, 7496 vec![Tag::from_parts("h", &[group_id]).expect("h")], 7497 content, 7498 ) 7499 } 7500 7501 fn runtime_pocket_event( 7502 key: FixtureKey, 7503 created_at: u64, 7504 kind: u64, 7505 tags: Vec<Tag>, 7506 content: &str, 7507 ) -> PocketOwnedEvent { 7508 let tags = pocket_tags_from_protocol(&tags); 7509 signed_pocket_event( 7510 fixture_secret_byte(key), 7511 created_at, 7512 u16::try_from(kind).expect("pocket kind"), 7513 &tags, 7514 content.as_bytes(), 7515 ) 7516 } 7517 7518 fn runtime_pocket_event_to_protocol(event: &PocketEvent) -> Result<Event, String> { 7519 let tags = event 7520 .tags() 7521 .map_err(|error| error.to_string())? 7522 .iter() 7523 .map(|tag| { 7524 Tag::new( 7525 tag.map(|value| { 7526 std::str::from_utf8(value) 7527 .map(str::to_owned) 7528 .map_err(|error| error.to_string()) 7529 }) 7530 .collect::<Result<Vec<_>, _>>()?, 7531 ) 7532 .map_err(|error| error.to_string()) 7533 }) 7534 .collect::<Result<Vec<_>, _>>()?; 7535 Ok(Event::new( 7536 EventId::new(&event.id().as_hex_string()).map_err(|error| error.to_string())?, 7537 UnsignedEvent::new( 7538 PublicKeyHex::new(&event.pubkey().as_hex_string()) 7539 .map_err(|error| error.to_string())?, 7540 UnixTimestamp::new(event.created_at().as_u64()), 7541 Kind::new(u64::from(event.kind().as_u16())).map_err(|error| error.to_string())?, 7542 tags, 7543 std::str::from_utf8(event.content()).map_err(|error| error.to_string())?, 7544 ), 7545 SignatureHex::new(&event.sig().to_string()).map_err(|error| error.to_string())?, 7546 )) 7547 } 7548 7549 fn pocket_tags_from_protocol(tags: &[Tag]) -> PocketOwnedTags { 7550 let parts = tags 7551 .iter() 7552 .map(|tag| tag.values().iter().map(String::as_str).collect::<Vec<_>>()) 7553 .collect::<Vec<_>>(); 7554 PocketOwnedTags::new(&parts).expect("pocket tags") 7555 } 7556 7557 fn fixture_secret_byte(key: FixtureKey) -> u8 { 7558 match key { 7559 FixtureKey::Relay => 9, 7560 FixtureKey::Admin => 11, 7561 FixtureKey::Member => 12, 7562 FixtureKey::Outsider => 13, 7563 FixtureKey::Owner => 10, 7564 } 7565 } 7566 7567 fn signed_pocket_event( 7568 secret_byte: u8, 7569 created_at: u64, 7570 kind: u16, 7571 tags: &PocketOwnedTags, 7572 content: &[u8], 7573 ) -> PocketOwnedEvent { 7574 let secret = format!("{secret_byte:02x}").repeat(32); 7575 RelaySigner::from_secret_hex(&secret) 7576 .expect("signer") 7577 .sign_pocket_event( 7578 PocketKind::from_u16(kind), 7579 tags, 7580 PocketTime::from_u64(created_at), 7581 content, 7582 ) 7583 .expect("pocket event") 7584 } 7585 7586 fn temp_root(name: &str) -> PathBuf { 7587 std::env::temp_dir().join(format!("tangle-runtime-{name}-{}", std::process::id())) 7588 } 7589 }