core.rs (176496B)
1 use crate::errors::{BaseRelayError, ok_accepted, ok_rejected}; 2 use crate::groups::{ 3 GroupEventWrite, GroupEventWriteError, GroupProjectionReadGuard, GroupServiceHandle, 4 }; 5 use crate::logging::{self, TangleModerationAuditResult}; 6 use crate::ops::BaseRelayReadinessState; 7 #[cfg(test)] 8 use crate::pocket_conversion::{tangle_event_to_pocket, tangle_filter_to_pocket}; 9 use crate::pocket_event_validation::{ 10 is_pocket_nip70_protected_event, pocket_event_id as pocket_runtime_event_id, pocket_event_kind, 11 pocket_event_pubkey, validate_pocket_event_shape, verify_pocket_event_signature, 12 }; 13 #[cfg(test)] 14 use crate::relay::outbound::protocol_messages_for_test; 15 use crate::relay::{ 16 auth::BaseAuthState, 17 filter::BaseRelayMatchedFilterContext, 18 live::{CloseResult, LiveSubscriptionSet}, 19 outbound::RuntimeRelayMessage, 20 }; 21 use std::{ 22 cell::{Cell, RefCell}, 23 collections::BTreeSet, 24 }; 25 use tangle_groups::{ 26 GroupAuthContext, GroupEventClass, GroupEventView, GroupId, GroupRuntimeConfig, 27 NIP29_RELAY_GENERATED_KIND_VALUES, StoreOffset, classify_group_event, 28 validate_client_group_event_structure, 29 }; 30 #[cfg(test)] 31 use tangle_protocol::{ClientMessage, Event, Filter}; 32 use tangle_protocol::{RelayMessage, SubscriptionId, UnixTimestamp}; 33 use tangle_store_pocket::{ 34 PocketEvent, PocketFilter, PocketHll8, PocketOwnedEvent, PocketOwnedFilter, PocketQueryConfig, 35 PocketScreenResult, PocketStoreConfig, PocketStoreHandle, 36 }; 37 38 pub(crate) const NEGENTROPY_DISABLED_MESSAGE: &str = "blocked: Negentropy sync is disabled"; 39 40 pub struct BaseRelay { 41 store: PocketStoreHandle, 42 subscriptions: LiveSubscriptionSet, 43 groups: Option<GroupServiceHandle>, 44 readiness: BaseRelayReadinessState, 45 limits: BaseRelayLimits, 46 query: PocketQueryConfig, 47 } 48 49 #[derive(Debug, Clone, PartialEq)] 50 pub(crate) struct BaseRelayEventWrite { 51 message: RelayMessage, 52 stored_offsets: Vec<StoreOffset>, 53 } 54 55 impl BaseRelayEventWrite { 56 fn stored(message: RelayMessage, stored_offsets: Vec<StoreOffset>) -> Self { 57 Self { 58 message, 59 stored_offsets, 60 } 61 } 62 63 fn unstored(message: RelayMessage) -> Self { 64 Self { 65 message, 66 stored_offsets: Vec::new(), 67 } 68 } 69 70 pub(crate) fn stored_offsets(&self) -> &[StoreOffset] { 71 &self.stored_offsets 72 } 73 74 pub(crate) fn into_message(self) -> RelayMessage { 75 self.message 76 } 77 } 78 79 #[derive(Debug, Clone, PartialEq)] 80 pub(crate) struct BaseRelayQueryReport { 81 messages: Vec<RuntimeRelayMessage>, 82 group_read_denied: bool, 83 query_metrics: BaseRelayQueryMetrics, 84 } 85 86 pub(crate) struct BaseRelayReqQuery<'a> { 87 subscription_id: SubscriptionId, 88 filters: Vec<PocketOwnedFilter>, 89 search_present: bool, 90 auth: &'a BaseAuthState, 91 } 92 93 impl<'a> BaseRelayReqQuery<'a> { 94 pub(crate) fn new( 95 subscription_id: SubscriptionId, 96 filters: Vec<PocketOwnedFilter>, 97 search_present: bool, 98 auth: &'a BaseAuthState, 99 ) -> Self { 100 Self { 101 subscription_id, 102 filters, 103 search_present, 104 auth, 105 } 106 } 107 } 108 109 struct BaseRelayGroupReqQuery<'a> { 110 subscription_id: SubscriptionId, 111 filters: Vec<PocketOwnedFilter>, 112 search_present: bool, 113 auth: &'a GroupAuthContext, 114 } 115 116 pub(crate) struct BaseRelayCountQuery<'a> { 117 subscription_id: SubscriptionId, 118 filters: Vec<PocketOwnedFilter>, 119 search_present: bool, 120 auth: &'a BaseAuthState, 121 } 122 123 impl<'a> BaseRelayCountQuery<'a> { 124 pub(crate) fn new( 125 subscription_id: SubscriptionId, 126 filters: Vec<PocketOwnedFilter>, 127 search_present: bool, 128 auth: &'a BaseAuthState, 129 ) -> Self { 130 Self { 131 subscription_id, 132 filters, 133 search_present, 134 auth, 135 } 136 } 137 } 138 139 struct BaseRelayGroupCountQuery<'a> { 140 subscription_id: SubscriptionId, 141 filters: Vec<PocketOwnedFilter>, 142 search_present: bool, 143 auth: &'a GroupAuthContext, 144 } 145 146 impl BaseRelayQueryReport { 147 pub(crate) fn new( 148 messages: Vec<RuntimeRelayMessage>, 149 group_read_denied: bool, 150 query_metrics: BaseRelayQueryMetrics, 151 ) -> Self { 152 Self { 153 messages, 154 group_read_denied, 155 query_metrics, 156 } 157 } 158 159 pub(crate) fn group_read_denied(&self) -> bool { 160 self.group_read_denied 161 } 162 163 pub(crate) fn query_metrics(&self) -> BaseRelayQueryMetrics { 164 self.query_metrics 165 } 166 167 pub(crate) fn into_messages(self) -> Vec<RuntimeRelayMessage> { 168 self.messages 169 } 170 } 171 172 #[derive(Debug, Clone, PartialEq)] 173 pub(crate) struct BaseRelayCountReport { 174 message: RelayMessage, 175 group_read_denied: bool, 176 query_metrics: BaseRelayQueryMetrics, 177 } 178 179 impl BaseRelayCountReport { 180 fn new( 181 message: RelayMessage, 182 group_read_denied: bool, 183 query_metrics: BaseRelayQueryMetrics, 184 ) -> Self { 185 Self { 186 message, 187 group_read_denied, 188 query_metrics, 189 } 190 } 191 192 pub(crate) fn group_read_denied(&self) -> bool { 193 self.group_read_denied 194 } 195 196 pub(crate) fn query_metrics(&self) -> BaseRelayQueryMetrics { 197 self.query_metrics 198 } 199 200 pub(crate) fn into_message(self) -> RelayMessage { 201 self.message 202 } 203 } 204 205 #[derive(Debug, Clone, PartialEq)] 206 pub(crate) struct BaseRelayEventQueryReport { 207 events: Vec<PocketOwnedEvent>, 208 group_read_denied: bool, 209 query_metrics: BaseRelayQueryMetrics, 210 } 211 212 impl BaseRelayEventQueryReport { 213 fn new( 214 events: Vec<PocketOwnedEvent>, 215 group_read_denied: bool, 216 query_metrics: BaseRelayQueryMetrics, 217 ) -> Self { 218 Self { 219 events, 220 group_read_denied, 221 query_metrics, 222 } 223 } 224 225 pub(crate) fn group_read_denied(&self) -> bool { 226 self.group_read_denied 227 } 228 229 pub(crate) fn query_metrics(&self) -> BaseRelayQueryMetrics { 230 self.query_metrics 231 } 232 233 pub(crate) fn into_events(self) -> Vec<PocketOwnedEvent> { 234 self.events 235 } 236 } 237 238 #[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] 239 pub(crate) struct BaseRelayQueryMetrics { 240 candidates_scanned: u64, 241 returned_events: u64, 242 redacted_events: u64, 243 } 244 245 impl BaseRelayQueryMetrics { 246 pub(crate) fn new(candidates_scanned: u64, returned_events: u64, redacted_events: u64) -> Self { 247 Self { 248 candidates_scanned, 249 returned_events, 250 redacted_events, 251 } 252 } 253 254 pub(crate) fn add(self, other: Self) -> Self { 255 Self { 256 candidates_scanned: self 257 .candidates_scanned 258 .saturating_add(other.candidates_scanned), 259 returned_events: self.returned_events.saturating_add(other.returned_events), 260 redacted_events: self.redacted_events.saturating_add(other.redacted_events), 261 } 262 } 263 264 pub(crate) fn with_returned_events(self, returned_events: usize) -> Self { 265 Self { 266 returned_events: u64::try_from(returned_events).expect("returned events fit in u64"), 267 ..self 268 } 269 } 270 271 pub(crate) fn candidates_scanned(self) -> u64 { 272 self.candidates_scanned 273 } 274 275 pub(crate) fn returned_events(self) -> u64 { 276 self.returned_events 277 } 278 279 pub(crate) fn redacted_events(self) -> u64 { 280 self.redacted_events 281 } 282 } 283 284 #[derive(Debug, Clone, PartialEq, Eq)] 285 struct BaseRelayCountEventsReport { 286 count: u64, 287 hll: Option<String>, 288 group_read_denied: bool, 289 query_metrics: BaseRelayQueryMetrics, 290 } 291 292 impl BaseRelayCountEventsReport { 293 fn new( 294 count: u64, 295 hll: Option<String>, 296 group_read_denied: bool, 297 query_metrics: BaseRelayQueryMetrics, 298 ) -> Self { 299 Self { 300 count, 301 hll, 302 group_read_denied, 303 query_metrics, 304 } 305 } 306 } 307 308 struct BaseRelayCountHll { 309 offset: Option<usize>, 310 hll: Option<PocketHll8>, 311 suppressed: bool, 312 } 313 314 #[derive(Debug, Clone, PartialEq, Eq)] 315 enum BaseRelayCountHllGroupTargets { 316 None, 317 Suppress, 318 Targets(Vec<GroupId>), 319 } 320 321 #[derive(Debug, Clone, Copy, PartialEq, Eq)] 322 enum BaseRelayCountHllTargetPolicy { 323 Eligible, 324 Suppress, 325 } 326 327 #[derive(Debug, Clone, Copy, PartialEq, Eq)] 328 enum BaseRelayCountHllDTagMode { 329 Ignore, 330 Target, 331 Suppress, 332 } 333 334 impl BaseRelayCountHll { 335 fn new(filters: &[PocketOwnedFilter]) -> Result<Self, BaseRelayError> { 336 let offset = BaseRelay::count_hll_offset(filters)?; 337 Ok(Self { 338 offset, 339 hll: offset.map(|_| PocketHll8::new()), 340 suppressed: false, 341 }) 342 } 343 344 fn suppress(&mut self) { 345 if self.offset.is_some() { 346 self.suppressed = true; 347 } 348 } 349 350 fn suppress_for_filter_targets( 351 &mut self, 352 groups: Option<&GroupServiceHandle>, 353 filters: &[PocketOwnedFilter], 354 ) { 355 if self.offset.is_none() { 356 return; 357 } 358 let [filter] = filters else { 359 return; 360 }; 361 if BaseRelay::count_hll_filter_target_policy(groups, filter) 362 == BaseRelayCountHllTargetPolicy::Suppress 363 { 364 self.suppress(); 365 } 366 } 367 368 fn observe( 369 &mut self, 370 groups: Option<&GroupServiceHandle>, 371 event: &PocketEvent, 372 ) -> Result<(), BaseRelayError> { 373 let Some(offset) = self.offset else { 374 return Ok(()); 375 }; 376 if BaseRelay::event_suppresses_count_hll(groups, event)? { 377 self.suppressed = true; 378 return Ok(()); 379 } 380 if let Some(hll) = &mut self.hll { 381 hll.add_element(event.pubkey().as_bytes(), offset) 382 .map_err(|error| BaseRelayError::error(error.to_string()))?; 383 } 384 Ok(()) 385 } 386 387 fn into_hex(self) -> Option<String> { 388 (!self.suppressed) 389 .then(|| self.hll.map(|value| value.to_hex_string())) 390 .flatten() 391 } 392 } 393 394 #[derive(Debug, Clone, Copy, PartialEq, Eq)] 395 pub(crate) enum BaseRelayFilterLimitMode { 396 ApplyDefaultLimit, 397 PreserveCountLimitless, 398 Override(u32), 399 } 400 401 #[derive(Debug, Clone, Copy, PartialEq, Eq)] 402 pub struct BaseRelayShutdownReport { 403 closed_subscriptions: usize, 404 } 405 406 impl BaseRelayShutdownReport { 407 pub fn new(closed_subscriptions: usize) -> Self { 408 Self { 409 closed_subscriptions, 410 } 411 } 412 413 pub fn closed_subscriptions(self) -> usize { 414 self.closed_subscriptions 415 } 416 } 417 418 #[derive(Debug, Clone, Copy, PartialEq, Eq)] 419 pub struct BaseRelayLimits { 420 max_pending_events: usize, 421 max_subscription_id_length: usize, 422 max_subscriptions: usize, 423 max_filters_per_request: usize, 424 max_tag_values_per_filter: usize, 425 max_query_complexity: usize, 426 max_event_tags: usize, 427 max_content_length: usize, 428 max_limit: u64, 429 default_limit: u64, 430 } 431 432 #[derive(Debug, Clone, Copy, PartialEq, Eq)] 433 pub struct BaseRelayLimitSettings { 434 pub max_pending_events: usize, 435 pub max_subscription_id_length: usize, 436 pub max_subscriptions: usize, 437 pub max_filters_per_request: usize, 438 pub max_tag_values_per_filter: usize, 439 pub max_query_complexity: usize, 440 pub max_event_tags: usize, 441 pub max_content_length: usize, 442 pub max_limit: u64, 443 pub default_limit: u64, 444 } 445 446 impl BaseRelayLimits { 447 pub fn new(settings: BaseRelayLimitSettings) -> Result<Self, BaseRelayError> { 448 let max_pending_events = settings.max_pending_events; 449 let max_subscription_id_length = settings.max_subscription_id_length; 450 let max_subscriptions = settings.max_subscriptions; 451 let max_filters_per_request = settings.max_filters_per_request; 452 let max_tag_values_per_filter = settings.max_tag_values_per_filter; 453 let max_query_complexity = settings.max_query_complexity; 454 let max_event_tags = settings.max_event_tags; 455 let max_content_length = settings.max_content_length; 456 let max_limit = settings.max_limit; 457 let default_limit = settings.default_limit; 458 if max_pending_events == 0 { 459 return Err(BaseRelayError::invalid( 460 "runtime max pending events must be greater than zero", 461 )); 462 } 463 if max_subscription_id_length == 0 { 464 return Err(BaseRelayError::invalid( 465 "runtime max subscription id length must be greater than zero", 466 )); 467 } 468 if max_subscriptions == 0 { 469 return Err(BaseRelayError::invalid( 470 "runtime max subscriptions per connection must be greater than zero", 471 )); 472 } 473 if max_filters_per_request == 0 { 474 return Err(BaseRelayError::invalid( 475 "runtime max filters per request must be greater than zero", 476 )); 477 } 478 if max_tag_values_per_filter == 0 { 479 return Err(BaseRelayError::invalid( 480 "runtime max tag values per filter must be greater than zero", 481 )); 482 } 483 if max_query_complexity == 0 { 484 return Err(BaseRelayError::invalid( 485 "runtime max query complexity must be greater than zero", 486 )); 487 } 488 if max_event_tags == 0 { 489 return Err(BaseRelayError::invalid( 490 "runtime max event tags must be greater than zero", 491 )); 492 } 493 if max_content_length == 0 { 494 return Err(BaseRelayError::invalid( 495 "runtime max content length must be greater than zero", 496 )); 497 } 498 if max_limit == 0 { 499 return Err(BaseRelayError::invalid( 500 "runtime max filter limit must be greater than zero", 501 )); 502 } 503 if default_limit == 0 { 504 return Err(BaseRelayError::invalid( 505 "runtime default filter limit must be greater than zero", 506 )); 507 } 508 if default_limit > max_limit { 509 return Err(BaseRelayError::invalid( 510 "runtime default filter limit must not exceed max filter limit", 511 )); 512 } 513 if usize::try_from(default_limit).is_ok_and(|limit| limit > max_query_complexity) { 514 return Err(BaseRelayError::invalid( 515 "runtime default filter limit must not exceed max query complexity", 516 )); 517 } 518 Ok(Self { 519 max_pending_events, 520 max_subscription_id_length, 521 max_subscriptions, 522 max_filters_per_request, 523 max_tag_values_per_filter, 524 max_query_complexity, 525 max_event_tags, 526 max_content_length, 527 max_limit, 528 default_limit, 529 }) 530 } 531 532 pub fn max_pending_events(self) -> usize { 533 self.max_pending_events 534 } 535 536 pub fn max_subscription_id_length(self) -> usize { 537 self.max_subscription_id_length 538 } 539 540 pub fn max_subscriptions(self) -> usize { 541 self.max_subscriptions 542 } 543 544 pub fn max_filters_per_request(self) -> usize { 545 self.max_filters_per_request 546 } 547 548 pub fn max_tag_values_per_filter(self) -> usize { 549 self.max_tag_values_per_filter 550 } 551 552 pub fn max_query_complexity(self) -> usize { 553 self.max_query_complexity 554 } 555 556 pub fn max_event_tags(self) -> usize { 557 self.max_event_tags 558 } 559 560 pub fn max_content_length(self) -> usize { 561 self.max_content_length 562 } 563 564 pub fn max_limit(self) -> u64 { 565 self.max_limit 566 } 567 568 pub fn default_limit(self) -> u64 { 569 self.default_limit 570 } 571 572 #[cfg(test)] 573 fn validate_protocol_event_for_test(&self, event: &Event) -> Result<(), BaseRelayError> { 574 if event.unsigned().tags().len() > self.max_event_tags { 575 return Err(BaseRelayError::invalid(format!( 576 "event tag count exceeds runtime max_event_tags {}", 577 self.max_event_tags 578 ))); 579 } 580 if event.unsigned().content().len() > self.max_content_length { 581 return Err(BaseRelayError::invalid(format!( 582 "event content length exceeds runtime max_content_length {}", 583 self.max_content_length 584 ))); 585 } 586 Ok(()) 587 } 588 589 pub(crate) fn validate_pocket_event(&self, event: &PocketEvent) -> Result<(), BaseRelayError> { 590 validate_pocket_event_shape(event, self.max_event_tags, self.max_content_length) 591 } 592 593 pub fn validate_subscription_id( 594 &self, 595 subscription_id: &SubscriptionId, 596 ) -> Result<(), BaseRelayError> { 597 let actual = subscription_id.as_str().chars().count(); 598 if actual > self.max_subscription_id_length { 599 return Err(BaseRelayError::invalid(format!( 600 "subscription id length exceeds runtime max_subid_length {}", 601 self.max_subscription_id_length 602 ))); 603 } 604 Ok(()) 605 } 606 607 pub(crate) fn validate_pocket_filters( 608 &self, 609 filters: &[PocketOwnedFilter], 610 ) -> Result<(), BaseRelayError> { 611 if filters.is_empty() { 612 return Err(BaseRelayError::invalid( 613 "request must include at least one filter", 614 )); 615 } 616 if filters.len() > self.max_filters_per_request { 617 return Err(BaseRelayError::invalid(format!( 618 "filter count exceeds runtime max_filters_per_request {}", 619 self.max_filters_per_request 620 ))); 621 } 622 for filter in filters { 623 let tag_values = filter 624 .tags() 625 .map_err(|error| BaseRelayError::error(error.to_string()))? 626 .iter() 627 .map(|tag| tag.skip(1).count()) 628 .sum::<usize>(); 629 if tag_values > self.max_tag_values_per_filter { 630 return Err(BaseRelayError::invalid(format!( 631 "filter tag value count exceeds runtime max_tag_values_per_filter {}", 632 self.max_tag_values_per_filter 633 ))); 634 } 635 if filter.limit() != u32::MAX && u64::from(filter.limit()) > self.max_limit { 636 return Err(BaseRelayError::invalid(format!( 637 "filter limit exceeds runtime max_limit {}", 638 self.max_limit 639 ))); 640 } 641 } 642 self.validate_pocket_query_complexity(filters)?; 643 Ok(()) 644 } 645 646 fn effective_pocket_filter_limit(self, filter: &PocketFilter) -> usize { 647 if filter.limit() == u32::MAX { 648 usize::try_from(self.default_limit).unwrap_or(usize::MAX) 649 } else { 650 usize::try_from(filter.limit()).unwrap_or(usize::MAX) 651 } 652 } 653 654 pub(crate) fn effective_pocket_filter_limit_for_query(self, filter: &PocketFilter) -> usize { 655 self.effective_pocket_filter_limit(filter) 656 } 657 658 fn validate_pocket_query_complexity( 659 &self, 660 filters: &[PocketOwnedFilter], 661 ) -> Result<(), BaseRelayError> { 662 let score = filters 663 .iter() 664 .map(|filter| self.pocket_filter_complexity(filter)) 665 .fold(0_usize, usize::saturating_add); 666 if score > self.max_query_complexity { 667 return Err(BaseRelayError::invalid(format!( 668 "query complexity {score} exceeds runtime max_query_complexity {}", 669 self.max_query_complexity 670 ))); 671 } 672 Ok(()) 673 } 674 675 fn pocket_filter_complexity(&self, filter: &PocketFilter) -> usize { 676 let tag_score = filter 677 .tags() 678 .map(|tags| { 679 tags.iter() 680 .map(|tag| 1_usize.saturating_add(tag.skip(1).count())) 681 .fold(0_usize, usize::saturating_add) 682 }) 683 .unwrap_or(usize::MAX); 684 1_usize 685 .saturating_add(filter.num_ids()) 686 .saturating_add(filter.num_authors()) 687 .saturating_add(filter.num_kinds()) 688 .saturating_add(tag_score) 689 .saturating_add(usize::from( 690 filter.since() != tangle_store_pocket::PocketTime::min(), 691 )) 692 .saturating_add(usize::from( 693 filter.until() != tangle_store_pocket::PocketTime::max(), 694 )) 695 .saturating_add(self.effective_pocket_filter_limit(filter)) 696 } 697 } 698 699 impl BaseRelay { 700 pub(crate) fn unsupported_search_present_closed( 701 subscription_id: &SubscriptionId, 702 search_present: bool, 703 ) -> Option<RelayMessage> { 704 search_present.then(|| RelayMessage::Closed { 705 subscription_id: subscription_id.clone(), 706 message: "unsupported: search filters are not supported".to_owned(), 707 }) 708 } 709 710 pub(crate) fn redacted_req_closed( 711 subscription_id: SubscriptionId, 712 auth: &GroupAuthContext, 713 ) -> RelayMessage { 714 let message = if auth.authenticated_pubkeys().is_empty() { 715 BaseRelayError::auth_required("authentication required to read group events") 716 .prefixed_message() 717 } else { 718 BaseRelayError::restricted("group is unavailable").prefixed_message() 719 }; 720 RelayMessage::Closed { 721 subscription_id, 722 message, 723 } 724 } 725 726 pub fn open( 727 config: &PocketStoreConfig, 728 limits: BaseRelayLimits, 729 query: PocketQueryConfig, 730 ) -> Result<Self, BaseRelayError> { 731 let store = PocketStoreHandle::open(config).map_err(BaseRelayError::from)?; 732 Self::new(store, limits, query) 733 } 734 735 pub fn open_with_groups( 736 config: &PocketStoreConfig, 737 limits: BaseRelayLimits, 738 groups: &GroupRuntimeConfig, 739 query: PocketQueryConfig, 740 ) -> Result<Self, BaseRelayError> { 741 let store = PocketStoreHandle::open(config).map_err(BaseRelayError::from)?; 742 Self::new_with_groups(store, limits, groups, query) 743 } 744 745 pub fn new( 746 store: PocketStoreHandle, 747 limits: BaseRelayLimits, 748 query: PocketQueryConfig, 749 ) -> Result<Self, BaseRelayError> { 750 Self::new_with_groups(store, limits, &GroupRuntimeConfig::disabled(), query) 751 } 752 753 pub fn new_with_groups( 754 store: PocketStoreHandle, 755 limits: BaseRelayLimits, 756 groups: &GroupRuntimeConfig, 757 query: PocketQueryConfig, 758 ) -> Result<Self, BaseRelayError> { 759 let groups = GroupServiceHandle::from_config(&store, groups)?; 760 let subscriptions = 761 LiveSubscriptionSet::new(limits.max_pending_events(), limits.max_subscriptions())?; 762 let readiness = BaseRelayReadinessState::runtime_ready_before_bind(); 763 Ok(Self { 764 store, 765 subscriptions, 766 groups, 767 readiness, 768 limits, 769 query, 770 }) 771 } 772 773 #[cfg(test)] 774 pub fn handle_client_message( 775 &mut self, 776 message: ClientMessage, 777 auth: &mut BaseAuthState, 778 now: UnixTimestamp, 779 ) -> Result<Vec<RelayMessage>, BaseRelayError> { 780 match message { 781 ClientMessage::Event(event) => self 782 .handle_event_with_auth(event, auth) 783 .map(|message| vec![message]), 784 ClientMessage::Req { 785 subscription_id, 786 filters, 787 } => self.handle_protocol_req_with_auth_for_test(subscription_id, filters, auth), 788 ClientMessage::Count { 789 subscription_id, 790 filters, 791 } => { 792 let search_present = filters.iter().any(|filter| filter.search().is_some()); 793 let filters = filters 794 .iter() 795 .map(tangle_filter_to_pocket) 796 .collect::<Result<Vec<_>, _>>()?; 797 self.handle_count_with_group_auth_report( 798 subscription_id, 799 filters, 800 search_present, 801 &GroupAuthContext::new(auth.authenticated_pubkeys().iter().cloned()), 802 ) 803 .map(|report| vec![report.into_message()]) 804 } 805 ClientMessage::Close(subscription_id) => { 806 self.handle_close(&subscription_id); 807 Ok(Vec::new()) 808 } 809 ClientMessage::Auth(event) => Ok(self.handle_auth_message(event, auth, now)), 810 ClientMessage::NegOpen { 811 subscription_id, .. 812 } 813 | ClientMessage::NegMsg { 814 subscription_id, .. 815 } => Ok(vec![Self::disabled_negentropy_message(subscription_id)]), 816 ClientMessage::NegClose(_) => Ok(Vec::new()), 817 } 818 } 819 820 pub(crate) fn disabled_negentropy_message(subscription_id: SubscriptionId) -> RelayMessage { 821 RelayMessage::NegErr { 822 subscription_id, 823 message: NEGENTROPY_DISABLED_MESSAGE.to_owned(), 824 } 825 } 826 827 pub(crate) fn query_req_with_shared_services( 828 store: &PocketStoreHandle, 829 groups: Option<&GroupServiceHandle>, 830 limits: BaseRelayLimits, 831 query: PocketQueryConfig, 832 request: BaseRelayReqQuery<'_>, 833 ) -> Result<BaseRelayQueryReport, BaseRelayError> { 834 let group_auth = 835 GroupAuthContext::new(request.auth.authenticated_pubkeys().iter().cloned()); 836 Self::query_req_with_group_auth_shared_services( 837 store, 838 groups, 839 limits, 840 query, 841 BaseRelayGroupReqQuery { 842 subscription_id: request.subscription_id, 843 filters: request.filters, 844 search_present: request.search_present, 845 auth: &group_auth, 846 }, 847 ) 848 } 849 850 fn event_by_offset(&self, offset: StoreOffset) -> Result<PocketOwnedEvent, BaseRelayError> { 851 self.store 852 .event_by_offset(offset.as_u64()) 853 .map_err(BaseRelayError::from) 854 } 855 856 pub fn event_by_offset_with_auth( 857 &self, 858 offset: StoreOffset, 859 auth: &BaseAuthState, 860 ) -> Result<Option<PocketOwnedEvent>, BaseRelayError> { 861 let event = self.event_by_offset(offset)?; 862 if Self::group_read_gate_visible_to_auth( 863 self.groups.as_ref(), 864 &event, 865 &GroupAuthContext::new(auth.authenticated_pubkeys().iter().cloned()), 866 )? { 867 Ok(Some(event)) 868 } else { 869 Ok(None) 870 } 871 } 872 873 #[cfg(test)] 874 fn handle_auth_message( 875 &self, 876 event: Event, 877 auth: &mut BaseAuthState, 878 now: UnixTimestamp, 879 ) -> Vec<RelayMessage> { 880 Self::handle_auth_with_limits(self.limits, event, auth, now) 881 } 882 883 #[cfg(test)] 884 pub(crate) fn handle_auth_with_limits( 885 limits: BaseRelayLimits, 886 event: Event, 887 auth: &mut BaseAuthState, 888 now: UnixTimestamp, 889 ) -> Vec<RelayMessage> { 890 if let Err(error) = limits.validate_protocol_event_for_test(&event) { 891 return vec![RelayMessage::Ok { 892 event_id: event.id().clone(), 893 accepted: false, 894 message: error.prefixed_message(), 895 }]; 896 } 897 auth.authenticate(&event, now) 898 .map(|_| { 899 vec![RelayMessage::Ok { 900 event_id: event.id().clone(), 901 accepted: true, 902 message: String::new(), 903 }] 904 }) 905 .unwrap_or_else(|error| { 906 vec![RelayMessage::Ok { 907 event_id: event.id().clone(), 908 accepted: false, 909 message: error.prefixed_message(), 910 }] 911 }) 912 } 913 914 pub(crate) fn handle_pocket_auth_with_limits( 915 limits: BaseRelayLimits, 916 event: &PocketEvent, 917 auth: &mut BaseAuthState, 918 now: UnixTimestamp, 919 ) -> Vec<RelayMessage> { 920 let event_id = 921 pocket_runtime_event_id(event).expect("Pocket event id is valid hex by construction"); 922 if let Err(error) = limits.validate_pocket_event(event) { 923 return vec![RelayMessage::Ok { 924 event_id, 925 accepted: false, 926 message: error.prefixed_message(), 927 }]; 928 } 929 auth.authenticate_pocket(event, now) 930 .map(|_| { 931 vec![RelayMessage::Ok { 932 event_id: event_id.clone(), 933 accepted: true, 934 message: String::new(), 935 }] 936 }) 937 .unwrap_or_else(|error| { 938 vec![RelayMessage::Ok { 939 event_id, 940 accepted: false, 941 message: error.prefixed_message(), 942 }] 943 }) 944 } 945 946 #[cfg(test)] 947 pub fn handle_event(&self, event: Event) -> Result<RelayMessage, BaseRelayError> { 948 self.handle_event_with_group_auth(event, &GroupAuthContext::unauthenticated()) 949 .map(BaseRelayEventWrite::into_message) 950 } 951 952 #[cfg(test)] 953 pub fn handle_event_with_auth( 954 &self, 955 event: Event, 956 auth: &BaseAuthState, 957 ) -> Result<RelayMessage, BaseRelayError> { 958 self.handle_event_with_auth_report(event, auth) 959 .map(BaseRelayEventWrite::into_message) 960 } 961 962 pub fn handle_pocket_event(&self, event: &PocketEvent) -> Result<RelayMessage, BaseRelayError> { 963 self.handle_pocket_event_with_group_auth(event, &GroupAuthContext::unauthenticated()) 964 .map(BaseRelayEventWrite::into_message) 965 } 966 967 pub fn handle_pocket_event_with_auth( 968 &self, 969 event: &PocketEvent, 970 auth: &BaseAuthState, 971 ) -> Result<RelayMessage, BaseRelayError> { 972 self.handle_pocket_event_with_auth_report(event, auth) 973 .map(BaseRelayEventWrite::into_message) 974 } 975 976 #[cfg(test)] 977 pub(crate) fn handle_event_with_auth_report( 978 &self, 979 event: Event, 980 auth: &BaseAuthState, 981 ) -> Result<BaseRelayEventWrite, BaseRelayError> { 982 Self::handle_event_with_shared_services( 983 &self.store, 984 self.groups.as_ref(), 985 self.limits, 986 event, 987 auth, 988 ) 989 } 990 991 #[cfg(test)] 992 pub(crate) fn handle_event_with_shared_services( 993 store: &PocketStoreHandle, 994 groups: Option<&GroupServiceHandle>, 995 limits: BaseRelayLimits, 996 event: Event, 997 auth: &BaseAuthState, 998 ) -> Result<BaseRelayEventWrite, BaseRelayError> { 999 Self::handle_event_with_group_auth_and_services( 1000 store, 1001 groups, 1002 limits, 1003 event, 1004 &GroupAuthContext::new(auth.authenticated_pubkeys().iter().cloned()), 1005 ) 1006 } 1007 1008 pub(crate) fn handle_pocket_event_with_auth_report( 1009 &self, 1010 event: &PocketEvent, 1011 auth: &BaseAuthState, 1012 ) -> Result<BaseRelayEventWrite, BaseRelayError> { 1013 Self::handle_pocket_event_with_shared_services( 1014 &self.store, 1015 self.groups.as_ref(), 1016 self.limits, 1017 event, 1018 auth, 1019 ) 1020 } 1021 1022 pub(crate) fn handle_pocket_event_with_shared_services( 1023 store: &PocketStoreHandle, 1024 groups: Option<&GroupServiceHandle>, 1025 limits: BaseRelayLimits, 1026 event: &PocketEvent, 1027 auth: &BaseAuthState, 1028 ) -> Result<BaseRelayEventWrite, BaseRelayError> { 1029 Self::handle_pocket_event_with_group_auth_and_services( 1030 store, 1031 groups, 1032 limits, 1033 event, 1034 &GroupAuthContext::new(auth.authenticated_pubkeys().iter().cloned()), 1035 ) 1036 } 1037 1038 pub fn groups_enabled(&self) -> bool { 1039 self.groups.is_some() 1040 } 1041 1042 pub(crate) fn store_handle(&self) -> PocketStoreHandle { 1043 self.store.clone() 1044 } 1045 1046 pub fn group_projection(&self) -> Option<GroupProjectionReadGuard<'_>> { 1047 self.groups.as_ref().map(GroupServiceHandle::projection) 1048 } 1049 1050 pub(crate) fn group_service_handle(&self) -> Option<GroupServiceHandle> { 1051 self.groups.clone() 1052 } 1053 1054 pub(crate) fn group_outbox_pending_events(&self) -> usize { 1055 self.groups 1056 .as_ref() 1057 .map(GroupServiceHandle::outbox_pending_events) 1058 .unwrap_or(0) 1059 } 1060 1061 pub fn readiness_state(&self) -> BaseRelayReadinessState { 1062 self.readiness.clone() 1063 } 1064 1065 pub fn shutdown(&mut self) -> Result<BaseRelayShutdownReport, BaseRelayError> { 1066 let closed = self.subscriptions.close_all(); 1067 self.store.sync()?; 1068 Ok(BaseRelayShutdownReport::new(closed)) 1069 } 1070 1071 #[cfg(test)] 1072 fn handle_event_with_group_auth( 1073 &self, 1074 event: Event, 1075 auth: &GroupAuthContext, 1076 ) -> Result<BaseRelayEventWrite, BaseRelayError> { 1077 Self::handle_event_with_group_auth_and_services( 1078 &self.store, 1079 self.groups.as_ref(), 1080 self.limits, 1081 event, 1082 auth, 1083 ) 1084 } 1085 1086 #[cfg(test)] 1087 fn handle_event_with_group_auth_and_services( 1088 store: &PocketStoreHandle, 1089 groups: Option<&GroupServiceHandle>, 1090 limits: BaseRelayLimits, 1091 event: Event, 1092 auth: &GroupAuthContext, 1093 ) -> Result<BaseRelayEventWrite, BaseRelayError> { 1094 let pocket_event = tangle_event_to_pocket(&event)?; 1095 Self::handle_pocket_event_with_group_auth_and_services( 1096 store, 1097 groups, 1098 limits, 1099 &pocket_event, 1100 auth, 1101 ) 1102 } 1103 1104 fn handle_pocket_event_with_group_auth( 1105 &self, 1106 event: &PocketEvent, 1107 auth: &GroupAuthContext, 1108 ) -> Result<BaseRelayEventWrite, BaseRelayError> { 1109 Self::handle_pocket_event_with_group_auth_and_services( 1110 &self.store, 1111 self.groups.as_ref(), 1112 self.limits, 1113 event, 1114 auth, 1115 ) 1116 } 1117 1118 fn handle_pocket_event_with_group_auth_and_services( 1119 store: &PocketStoreHandle, 1120 groups: Option<&GroupServiceHandle>, 1121 limits: BaseRelayLimits, 1122 event: &PocketEvent, 1123 auth: &GroupAuthContext, 1124 ) -> Result<BaseRelayEventWrite, BaseRelayError> { 1125 let event_id = pocket_runtime_event_id(event)?; 1126 if let Err(error) = limits.validate_pocket_event(event) { 1127 return Ok(BaseRelayEventWrite::unstored(ok_rejected( 1128 event_id, 1129 error.prefixed_message(), 1130 ))); 1131 } 1132 if let Err(error) = verify_pocket_event_signature(event) { 1133 return Ok(BaseRelayEventWrite::unstored(ok_rejected( 1134 event_id, 1135 error.prefixed_message(), 1136 ))); 1137 } 1138 let pubkey = pocket_event_pubkey(event)?; 1139 if is_pocket_nip70_protected_event(event)? && !auth.contains(&pubkey) { 1140 return Ok(BaseRelayEventWrite::unstored(ok_rejected( 1141 event_id, 1142 BaseRelayError::auth_required( 1143 "protected event requires authenticated event author", 1144 ) 1145 .prefixed_message(), 1146 ))); 1147 } 1148 let group_limits = groups.map(GroupServiceHandle::limits).unwrap_or_default(); 1149 let audit_class = classify_group_event(event, group_limits).ok(); 1150 let class = match validate_client_group_event_structure(event, group_limits) { 1151 Ok(class) => class, 1152 Err(error) => { 1153 if let Some(class) = audit_class.as_ref() { 1154 logging::log_group_moderation_audit( 1155 event, 1156 class, 1157 TangleModerationAuditResult::Rejected, 1158 ); 1159 } 1160 return Ok(BaseRelayEventWrite::unstored(ok_rejected( 1161 event_id, 1162 error.prefixed_message(), 1163 ))); 1164 } 1165 }; 1166 if !matches!(class, GroupEventClass::NonGroup) { 1167 let Some(groups) = groups else { 1168 logging::log_group_moderation_audit( 1169 event, 1170 &class, 1171 TangleModerationAuditResult::Rejected, 1172 ); 1173 return Ok(BaseRelayEventWrite::unstored(ok_rejected( 1174 event_id, 1175 "blocked: NIP-29 group events are not accepted before group service".to_owned(), 1176 ))); 1177 }; 1178 match groups.store_group_pocket_event(store, event, &class, auth) { 1179 Ok(GroupEventWrite::Stored(stored_offsets)) => { 1180 logging::log_group_moderation_audit( 1181 event, 1182 &class, 1183 TangleModerationAuditResult::Accepted, 1184 ); 1185 return Ok(BaseRelayEventWrite::stored( 1186 ok_accepted(event_id, String::new()), 1187 stored_offsets, 1188 )); 1189 } 1190 Ok(GroupEventWrite::Duplicate) => { 1191 logging::log_group_moderation_audit( 1192 event, 1193 &class, 1194 TangleModerationAuditResult::Accepted, 1195 ); 1196 return Ok(BaseRelayEventWrite::unstored(ok_accepted( 1197 event_id, 1198 "duplicate: already have this event".to_owned(), 1199 ))); 1200 } 1201 Err(GroupEventWriteError::Rejected(error)) => { 1202 logging::log_group_moderation_audit( 1203 event, 1204 &class, 1205 TangleModerationAuditResult::Rejected, 1206 ); 1207 return Ok(BaseRelayEventWrite::unstored(ok_rejected( 1208 event_id, 1209 error.prefixed_message(), 1210 ))); 1211 } 1212 Err(GroupEventWriteError::Storage(error)) => return Err(error), 1213 } 1214 } 1215 if pocket_event_kind(event)?.is_ephemeral() { 1216 return Ok(BaseRelayEventWrite::unstored(ok_accepted( 1217 event_id, 1218 String::new(), 1219 ))); 1220 } 1221 if store.event_by_id(event.id())?.is_some() { 1222 return Ok(BaseRelayEventWrite::unstored(ok_accepted( 1223 event_id, 1224 "duplicate: already have this event".to_owned(), 1225 ))); 1226 } 1227 let store_offset = StoreOffset::new(store.store_event(event)?); 1228 Ok(BaseRelayEventWrite::stored( 1229 ok_accepted(event_id, String::new()), 1230 vec![store_offset], 1231 )) 1232 } 1233 1234 pub fn handle_pocket_req( 1235 &mut self, 1236 subscription_id: SubscriptionId, 1237 filters: Vec<PocketOwnedFilter>, 1238 ) -> Result<Vec<RuntimeRelayMessage>, BaseRelayError> { 1239 self.handle_pocket_req_with_group_auth( 1240 subscription_id, 1241 filters, 1242 &GroupAuthContext::unauthenticated(), 1243 ) 1244 } 1245 1246 #[cfg(test)] 1247 pub fn handle_protocol_req_for_test( 1248 &mut self, 1249 subscription_id: SubscriptionId, 1250 filters: Vec<Filter>, 1251 ) -> Result<Vec<RelayMessage>, BaseRelayError> { 1252 self.handle_protocol_req_with_group_auth_for_test( 1253 subscription_id, 1254 filters, 1255 &GroupAuthContext::unauthenticated(), 1256 ) 1257 } 1258 1259 #[cfg(test)] 1260 pub fn handle_protocol_req_with_auth_for_test( 1261 &mut self, 1262 subscription_id: SubscriptionId, 1263 filters: Vec<Filter>, 1264 auth: &BaseAuthState, 1265 ) -> Result<Vec<RelayMessage>, BaseRelayError> { 1266 self.handle_protocol_req_with_group_auth_for_test( 1267 subscription_id, 1268 filters, 1269 &GroupAuthContext::new(auth.authenticated_pubkeys().iter().cloned()), 1270 ) 1271 } 1272 1273 #[cfg(test)] 1274 fn handle_protocol_req_with_group_auth_for_test( 1275 &mut self, 1276 subscription_id: SubscriptionId, 1277 filters: Vec<Filter>, 1278 auth: &GroupAuthContext, 1279 ) -> Result<Vec<RelayMessage>, BaseRelayError> { 1280 self.handle_protocol_req_with_group_auth_report_for_test(subscription_id, filters, auth) 1281 .map(BaseRelayQueryReport::into_messages) 1282 .and_then(protocol_messages_for_test) 1283 } 1284 1285 pub fn handle_pocket_req_with_auth( 1286 &mut self, 1287 subscription_id: SubscriptionId, 1288 filters: Vec<PocketOwnedFilter>, 1289 auth: &BaseAuthState, 1290 ) -> Result<Vec<RuntimeRelayMessage>, BaseRelayError> { 1291 self.handle_pocket_req_with_group_auth( 1292 subscription_id, 1293 filters, 1294 &GroupAuthContext::new(auth.authenticated_pubkeys().iter().cloned()), 1295 ) 1296 } 1297 1298 fn handle_pocket_req_with_group_auth( 1299 &mut self, 1300 subscription_id: SubscriptionId, 1301 filters: Vec<PocketOwnedFilter>, 1302 auth: &GroupAuthContext, 1303 ) -> Result<Vec<RuntimeRelayMessage>, BaseRelayError> { 1304 self.handle_pocket_req_with_group_auth_report(subscription_id, filters, false, auth) 1305 .map(BaseRelayQueryReport::into_messages) 1306 } 1307 1308 #[cfg(test)] 1309 fn handle_protocol_req_with_group_auth_report_for_test( 1310 &mut self, 1311 subscription_id: SubscriptionId, 1312 filters: Vec<Filter>, 1313 auth: &GroupAuthContext, 1314 ) -> Result<BaseRelayQueryReport, BaseRelayError> { 1315 let search_present = filters.iter().any(|filter| filter.search().is_some()); 1316 let filters = filters 1317 .iter() 1318 .map(tangle_filter_to_pocket) 1319 .collect::<Result<Vec<_>, _>>()?; 1320 self.handle_pocket_req_with_group_auth_report( 1321 subscription_id, 1322 filters, 1323 search_present, 1324 auth, 1325 ) 1326 } 1327 1328 fn handle_pocket_req_with_group_auth_report( 1329 &mut self, 1330 subscription_id: SubscriptionId, 1331 filters: Vec<PocketOwnedFilter>, 1332 search_present: bool, 1333 auth: &GroupAuthContext, 1334 ) -> Result<BaseRelayQueryReport, BaseRelayError> { 1335 self.limits.validate_subscription_id(&subscription_id)?; 1336 self.limits.validate_pocket_filters(&filters)?; 1337 if let Some(message) = 1338 Self::unsupported_search_present_closed(&subscription_id, search_present) 1339 { 1340 return Ok(BaseRelayQueryReport::new( 1341 vec![message.into()], 1342 false, 1343 BaseRelayQueryMetrics::default(), 1344 )); 1345 } 1346 let should_subscribe = !pocket_filters_are_complete(&filters); 1347 if should_subscribe { 1348 self.subscriptions 1349 .ensure_can_subscribe(&subscription_id, &filters)?; 1350 let report = self.query_req_with_group_auth_report( 1351 subscription_id.clone(), 1352 filters.clone(), 1353 false, 1354 auth, 1355 )?; 1356 if !report.group_read_denied() { 1357 self.subscriptions.subscribe(subscription_id, filters)?; 1358 } 1359 return Ok(report); 1360 } 1361 self.query_req_with_group_auth_report(subscription_id, filters, false, auth) 1362 } 1363 1364 fn query_req_with_group_auth_report( 1365 &self, 1366 subscription_id: SubscriptionId, 1367 filters: Vec<PocketOwnedFilter>, 1368 search_present: bool, 1369 auth: &GroupAuthContext, 1370 ) -> Result<BaseRelayQueryReport, BaseRelayError> { 1371 Self::query_req_with_group_auth_shared_services( 1372 &self.store, 1373 self.groups.as_ref(), 1374 self.limits, 1375 self.query, 1376 BaseRelayGroupReqQuery { 1377 subscription_id, 1378 filters, 1379 search_present, 1380 auth, 1381 }, 1382 ) 1383 } 1384 1385 fn query_req_with_group_auth_shared_services( 1386 store: &PocketStoreHandle, 1387 groups: Option<&GroupServiceHandle>, 1388 limits: BaseRelayLimits, 1389 query: PocketQueryConfig, 1390 request: BaseRelayGroupReqQuery<'_>, 1391 ) -> Result<BaseRelayQueryReport, BaseRelayError> { 1392 let BaseRelayGroupReqQuery { 1393 subscription_id, 1394 filters, 1395 search_present, 1396 auth, 1397 } = request; 1398 limits.validate_subscription_id(&subscription_id)?; 1399 limits.validate_pocket_filters(&filters)?; 1400 if let Some(message) = 1401 Self::unsupported_search_present_closed(&subscription_id, search_present) 1402 { 1403 return Ok(BaseRelayQueryReport::new( 1404 vec![message.into()], 1405 false, 1406 BaseRelayQueryMetrics::default(), 1407 )); 1408 } 1409 let report = 1410 Self::query_events_report_with_services(store, groups, limits, query, &filters, auth)?; 1411 let group_read_denied = report.group_read_denied; 1412 let query_metrics = report.query_metrics; 1413 let mut messages = report 1414 .events 1415 .into_iter() 1416 .map(|event| RuntimeRelayMessage::event(subscription_id.clone(), event)) 1417 .collect::<Vec<_>>(); 1418 if group_read_denied { 1419 messages.push(Self::redacted_req_closed(subscription_id, auth).into()); 1420 } else { 1421 messages.push(RelayMessage::Eose(subscription_id).into()); 1422 } 1423 Ok(BaseRelayQueryReport::new( 1424 messages, 1425 group_read_denied, 1426 query_metrics, 1427 )) 1428 } 1429 1430 pub fn handle_count( 1431 &self, 1432 subscription_id: SubscriptionId, 1433 filters: Vec<PocketOwnedFilter>, 1434 ) -> Result<RelayMessage, BaseRelayError> { 1435 self.handle_count_with_group_auth( 1436 subscription_id, 1437 filters, 1438 &GroupAuthContext::unauthenticated(), 1439 ) 1440 } 1441 1442 pub fn handle_count_with_auth( 1443 &self, 1444 subscription_id: SubscriptionId, 1445 filters: Vec<PocketOwnedFilter>, 1446 auth: &BaseAuthState, 1447 ) -> Result<RelayMessage, BaseRelayError> { 1448 self.handle_count_with_auth_report(subscription_id, filters, auth) 1449 .map(BaseRelayCountReport::into_message) 1450 } 1451 1452 pub(crate) fn handle_count_with_auth_report( 1453 &self, 1454 subscription_id: SubscriptionId, 1455 filters: Vec<PocketOwnedFilter>, 1456 auth: &BaseAuthState, 1457 ) -> Result<BaseRelayCountReport, BaseRelayError> { 1458 Self::handle_count_with_shared_services( 1459 &self.store, 1460 self.groups.as_ref(), 1461 self.limits, 1462 self.query, 1463 BaseRelayCountQuery::new(subscription_id, filters, false, auth), 1464 ) 1465 } 1466 1467 pub(crate) fn handle_count_with_shared_services( 1468 store: &PocketStoreHandle, 1469 groups: Option<&GroupServiceHandle>, 1470 limits: BaseRelayLimits, 1471 query: PocketQueryConfig, 1472 request: BaseRelayCountQuery<'_>, 1473 ) -> Result<BaseRelayCountReport, BaseRelayError> { 1474 let group_auth = 1475 GroupAuthContext::new(request.auth.authenticated_pubkeys().iter().cloned()); 1476 Self::handle_count_with_group_auth_shared_services( 1477 store, 1478 groups, 1479 limits, 1480 query, 1481 BaseRelayGroupCountQuery { 1482 subscription_id: request.subscription_id, 1483 filters: request.filters, 1484 search_present: request.search_present, 1485 auth: &group_auth, 1486 }, 1487 ) 1488 } 1489 1490 fn handle_count_with_group_auth( 1491 &self, 1492 subscription_id: SubscriptionId, 1493 filters: Vec<PocketOwnedFilter>, 1494 auth: &GroupAuthContext, 1495 ) -> Result<RelayMessage, BaseRelayError> { 1496 self.handle_count_with_group_auth_report(subscription_id, filters, false, auth) 1497 .map(BaseRelayCountReport::into_message) 1498 } 1499 1500 fn handle_count_with_group_auth_report( 1501 &self, 1502 subscription_id: SubscriptionId, 1503 filters: Vec<PocketOwnedFilter>, 1504 search_present: bool, 1505 auth: &GroupAuthContext, 1506 ) -> Result<BaseRelayCountReport, BaseRelayError> { 1507 Self::handle_count_with_group_auth_shared_services( 1508 &self.store, 1509 self.groups.as_ref(), 1510 self.limits, 1511 self.query, 1512 BaseRelayGroupCountQuery { 1513 subscription_id, 1514 filters, 1515 search_present, 1516 auth, 1517 }, 1518 ) 1519 } 1520 1521 fn handle_count_with_group_auth_shared_services( 1522 store: &PocketStoreHandle, 1523 groups: Option<&GroupServiceHandle>, 1524 limits: BaseRelayLimits, 1525 query: PocketQueryConfig, 1526 request: BaseRelayGroupCountQuery<'_>, 1527 ) -> Result<BaseRelayCountReport, BaseRelayError> { 1528 let BaseRelayGroupCountQuery { 1529 subscription_id, 1530 filters, 1531 search_present, 1532 auth, 1533 } = request; 1534 limits.validate_subscription_id(&subscription_id)?; 1535 limits.validate_pocket_filters(&filters)?; 1536 if let Some(message) = 1537 Self::unsupported_search_present_closed(&subscription_id, search_present) 1538 { 1539 return Ok(BaseRelayCountReport::new( 1540 message, 1541 false, 1542 BaseRelayQueryMetrics::default(), 1543 )); 1544 } 1545 let report = 1546 Self::count_events_report_with_services(store, groups, limits, query, &filters, auth)?; 1547 Ok(BaseRelayCountReport::new( 1548 RelayMessage::Count { 1549 subscription_id, 1550 count: report.count, 1551 hll: report.hll, 1552 }, 1553 report.group_read_denied, 1554 report.query_metrics, 1555 )) 1556 } 1557 1558 pub fn handle_close(&mut self, subscription_id: &SubscriptionId) -> CloseResult { 1559 self.subscriptions.close(subscription_id) 1560 } 1561 1562 pub fn fanout_pocket(&mut self, event: &PocketEvent) -> Vec<RuntimeRelayMessage> { 1563 self.fanout_pocket_with_group_auth(event, &GroupAuthContext::unauthenticated()) 1564 } 1565 1566 pub fn fanout_pocket_with_group_auth( 1567 &mut self, 1568 event: &PocketEvent, 1569 auth: &GroupAuthContext, 1570 ) -> Vec<RuntimeRelayMessage> { 1571 let groups = self.groups.as_ref(); 1572 self.subscriptions 1573 .fanout(event, auth, |event, auth| { 1574 Self::group_read_gate_visible_to_auth(groups, event, auth).unwrap_or(false) 1575 }) 1576 .expect("Pocket live fanout must match") 1577 .into_iter() 1578 .map(|matched| RuntimeRelayMessage::Event { 1579 subscription_id: matched.into_subscription_id(), 1580 event: event.to_owned(), 1581 }) 1582 .collect() 1583 } 1584 1585 #[cfg(test)] 1586 pub fn fanout_protocol_for_test(&mut self, event: &Event) -> Vec<RelayMessage> { 1587 self.fanout_protocol_with_group_auth_for_test(event, &GroupAuthContext::unauthenticated()) 1588 } 1589 1590 #[cfg(test)] 1591 pub fn fanout_protocol_with_group_auth_for_test( 1592 &mut self, 1593 event: &Event, 1594 auth: &GroupAuthContext, 1595 ) -> Vec<RelayMessage> { 1596 let pocket_event = tangle_event_to_pocket(event).expect("event must convert to Pocket"); 1597 protocol_messages_for_test(self.fanout_pocket_with_group_auth(&pocket_event, auth)) 1598 .expect("test protocol fanout must convert") 1599 } 1600 1601 pub fn active_subscription_count(&self) -> usize { 1602 self.subscriptions.active_count() 1603 } 1604 1605 fn query_events_report_with_services( 1606 store: &PocketStoreHandle, 1607 groups: Option<&GroupServiceHandle>, 1608 limits: BaseRelayLimits, 1609 query: PocketQueryConfig, 1610 filters: &[PocketOwnedFilter], 1611 auth: &GroupAuthContext, 1612 ) -> Result<BaseRelayEventQueryReport, BaseRelayError> { 1613 let mut output = Vec::new(); 1614 let mut group_read_denied = false; 1615 let mut query_metrics = BaseRelayQueryMetrics::default(); 1616 for filter in filters { 1617 let report = Self::query_filter_events_report_with_services( 1618 store, 1619 groups, 1620 limits, 1621 query, 1622 filter, 1623 auth, 1624 BaseRelayFilterLimitMode::ApplyDefaultLimit, 1625 )?; 1626 group_read_denied |= report.group_read_denied; 1627 query_metrics = query_metrics.add(report.query_metrics); 1628 let mut events = Self::sort_and_dedupe_query_events(report.events); 1629 events.truncate(limits.effective_pocket_filter_limit(filter)); 1630 output.extend(events); 1631 } 1632 let events = Self::sort_and_dedupe_query_events(output); 1633 query_metrics = query_metrics.with_returned_events(events.len()); 1634 Ok(BaseRelayEventQueryReport::new( 1635 events, 1636 group_read_denied, 1637 query_metrics, 1638 )) 1639 } 1640 1641 fn count_events_report_with_services( 1642 store: &PocketStoreHandle, 1643 groups: Option<&GroupServiceHandle>, 1644 limits: BaseRelayLimits, 1645 query: PocketQueryConfig, 1646 filters: &[PocketOwnedFilter], 1647 auth: &GroupAuthContext, 1648 ) -> Result<BaseRelayCountEventsReport, BaseRelayError> { 1649 let mut seen = BTreeSet::new(); 1650 let mut group_read_denied = false; 1651 let mut query_metrics = BaseRelayQueryMetrics::default(); 1652 let count_query = query.exact_count(); 1653 let mut hll = BaseRelayCountHll::new(filters)?; 1654 hll.suppress_for_filter_targets(groups, filters); 1655 for filter in filters { 1656 let report = Self::query_filter_events_report_with_services( 1657 store, 1658 groups, 1659 limits, 1660 count_query, 1661 filter, 1662 auth, 1663 BaseRelayFilterLimitMode::PreserveCountLimitless, 1664 )?; 1665 group_read_denied |= report.group_read_denied; 1666 if report.group_read_denied { 1667 hll.suppress(); 1668 } 1669 query_metrics = query_metrics.add(report.query_metrics); 1670 for event in report.events { 1671 let event: &PocketEvent = &event; 1672 hll.observe(groups, event)?; 1673 seen.insert(event.id()); 1674 } 1675 } 1676 let count = u64::try_from(seen.len()) 1677 .map_err(|_| BaseRelayError::error("visible event count overflow"))?; 1678 let hll = hll.into_hex(); 1679 Ok(BaseRelayCountEventsReport::new( 1680 count, 1681 hll, 1682 group_read_denied, 1683 query_metrics, 1684 )) 1685 } 1686 1687 fn count_hll_offset(filters: &[PocketOwnedFilter]) -> Result<Option<usize>, BaseRelayError> { 1688 let [filter] = filters else { 1689 return Ok(None); 1690 }; 1691 filter 1692 .hyperloglog_offset() 1693 .map_err(|error| BaseRelayError::error(error.to_string())) 1694 } 1695 1696 fn event_suppresses_count_hll( 1697 groups: Option<&GroupServiceHandle>, 1698 event: &PocketEvent, 1699 ) -> Result<bool, BaseRelayError> { 1700 let Some(groups) = groups else { 1701 return Ok(false); 1702 }; 1703 let class = classify_group_event(event, groups.limits()).map_err(BaseRelayError::from)?; 1704 let Some(group_id) = class.group_id() else { 1705 return Ok(false); 1706 }; 1707 let projection = groups.projection(); 1708 let Some(group) = projection.group(group_id) else { 1709 return Ok(true); 1710 }; 1711 Ok(projection.tombstone(group_id).is_some() 1712 || group.metadata().private() 1713 || group.metadata().hidden()) 1714 } 1715 1716 fn count_hll_filter_target_policy( 1717 groups: Option<&GroupServiceHandle>, 1718 filter: &PocketFilter, 1719 ) -> BaseRelayCountHllTargetPolicy { 1720 let Some(groups) = groups else { 1721 return if Self::count_hll_filter_has_group_target(filter) { 1722 BaseRelayCountHllTargetPolicy::Suppress 1723 } else { 1724 BaseRelayCountHllTargetPolicy::Eligible 1725 }; 1726 }; 1727 match Self::count_hll_group_targets( 1728 filter, 1729 usize::from(groups.limits().max_group_id_bytes()), 1730 ) { 1731 BaseRelayCountHllGroupTargets::None => BaseRelayCountHllTargetPolicy::Eligible, 1732 BaseRelayCountHllGroupTargets::Suppress => BaseRelayCountHllTargetPolicy::Suppress, 1733 BaseRelayCountHllGroupTargets::Targets(group_ids) => { 1734 let projection = groups.projection(); 1735 if group_ids.iter().all(|group_id| { 1736 projection.group(group_id).is_some_and(|group| { 1737 projection.tombstone(group_id).is_none() 1738 && !group.metadata().private() 1739 && !group.metadata().hidden() 1740 }) 1741 }) { 1742 BaseRelayCountHllTargetPolicy::Eligible 1743 } else { 1744 BaseRelayCountHllTargetPolicy::Suppress 1745 } 1746 } 1747 } 1748 } 1749 1750 fn count_hll_group_targets( 1751 filter: &PocketFilter, 1752 max_group_id_bytes: usize, 1753 ) -> BaseRelayCountHllGroupTargets { 1754 let Ok(tags) = filter.tags() else { 1755 return BaseRelayCountHllGroupTargets::Suppress; 1756 }; 1757 let d_tag_mode = Self::count_hll_filter_d_tag_mode(filter); 1758 let mut group_ids = Vec::new(); 1759 for tag in tags.iter() { 1760 let mut values = tag.into_iter(); 1761 let Some(name) = values.next() else { 1762 continue; 1763 }; 1764 if name == b"d" { 1765 match d_tag_mode { 1766 BaseRelayCountHllDTagMode::Ignore => continue, 1767 BaseRelayCountHllDTagMode::Suppress => { 1768 return BaseRelayCountHllGroupTargets::Suppress; 1769 } 1770 BaseRelayCountHllDTagMode::Target => {} 1771 } 1772 } else if name != b"h" { 1773 continue; 1774 } 1775 let mut found_value = false; 1776 for value in values { 1777 found_value = true; 1778 let Ok(value) = std::str::from_utf8(value) else { 1779 return BaseRelayCountHllGroupTargets::Suppress; 1780 }; 1781 let Ok(group_id) = GroupId::new_with_max_bytes(value, max_group_id_bytes) else { 1782 return BaseRelayCountHllGroupTargets::Suppress; 1783 }; 1784 group_ids.push(group_id); 1785 } 1786 if !found_value { 1787 return BaseRelayCountHllGroupTargets::Suppress; 1788 } 1789 } 1790 if group_ids.is_empty() { 1791 BaseRelayCountHllGroupTargets::None 1792 } else { 1793 group_ids.sort(); 1794 group_ids.dedup(); 1795 BaseRelayCountHllGroupTargets::Targets(group_ids) 1796 } 1797 } 1798 1799 fn count_hll_filter_has_group_target(filter: &PocketFilter) -> bool { 1800 let Ok(tags) = filter.tags() else { 1801 return true; 1802 }; 1803 let d_tag_mode = Self::count_hll_filter_d_tag_mode(filter); 1804 tags.iter().any(|tag| { 1805 let mut values = tag.into_iter(); 1806 let name = values.next(); 1807 matches!(name, Some(b"h")) 1808 || (matches!( 1809 d_tag_mode, 1810 BaseRelayCountHllDTagMode::Target | BaseRelayCountHllDTagMode::Suppress 1811 ) && matches!(name, Some(b"d"))) 1812 }) 1813 } 1814 1815 fn count_hll_filter_d_tag_mode(filter: &PocketFilter) -> BaseRelayCountHllDTagMode { 1816 if filter.num_kinds() == 0 { 1817 return BaseRelayCountHllDTagMode::Suppress; 1818 } 1819 if filter 1820 .kinds() 1821 .any(|kind| NIP29_RELAY_GENERATED_KIND_VALUES.contains(&u32::from(kind.as_u16()))) 1822 { 1823 BaseRelayCountHllDTagMode::Target 1824 } else { 1825 BaseRelayCountHllDTagMode::Ignore 1826 } 1827 } 1828 1829 pub(crate) fn query_filter_events_report_with_services( 1830 store: &PocketStoreHandle, 1831 groups: Option<&GroupServiceHandle>, 1832 limits: BaseRelayLimits, 1833 query: PocketQueryConfig, 1834 filter: &PocketFilter, 1835 auth: &GroupAuthContext, 1836 limit_mode: BaseRelayFilterLimitMode, 1837 ) -> Result<BaseRelayEventQueryReport, BaseRelayError> { 1838 let pocket_filter = Self::pocket_filter_with_limit_mode(limits, filter, limit_mode)?; 1839 let screen_error = RefCell::new(None); 1840 let candidates_scanned = Cell::new(0_u64); 1841 let redacted_events = Cell::new(0_u64); 1842 let screened = store.find_events_with_screen(&pocket_filter, query, |pocket_event| { 1843 candidates_scanned.set(candidates_scanned.get().saturating_add(1)); 1844 if screen_error.borrow().is_some() { 1845 return PocketScreenResult::Mismatch; 1846 } 1847 match pocket_filter.event_matches(pocket_event) { 1848 Ok(false) => PocketScreenResult::Mismatch, 1849 Ok(true) => { 1850 match Self::group_read_gate_visible_to_auth(groups, pocket_event, auth) { 1851 Ok(true) => PocketScreenResult::Match, 1852 Ok(false) => { 1853 redacted_events.set(redacted_events.get().saturating_add(1)); 1854 PocketScreenResult::Redacted 1855 } 1856 Err(error) => { 1857 *screen_error.borrow_mut() = Some(error); 1858 PocketScreenResult::Mismatch 1859 } 1860 } 1861 } 1862 Err(error) => { 1863 *screen_error.borrow_mut() = Some(BaseRelayError::error(error.to_string())); 1864 PocketScreenResult::Mismatch 1865 } 1866 } 1867 })?; 1868 if let Some(error) = screen_error.into_inner() { 1869 return Err(error); 1870 } 1871 let group_read_denied = screened.redacted(); 1872 let events = screened.into_events(); 1873 Ok(BaseRelayEventQueryReport::new( 1874 events, 1875 group_read_denied, 1876 BaseRelayQueryMetrics::new(candidates_scanned.get(), 0, redacted_events.get()), 1877 )) 1878 } 1879 1880 fn pocket_filter_with_limit_mode( 1881 limits: BaseRelayLimits, 1882 filter: &PocketFilter, 1883 limit_mode: BaseRelayFilterLimitMode, 1884 ) -> Result<PocketOwnedFilter, BaseRelayError> { 1885 let limit = match (limit_mode, filter.limit()) { 1886 (BaseRelayFilterLimitMode::ApplyDefaultLimit, u32::MAX) => { 1887 u32::try_from(limits.default_limit) 1888 .map_err(|_| BaseRelayError::invalid("default filter limit exceeds u32"))? 1889 } 1890 (BaseRelayFilterLimitMode::PreserveCountLimitless, _) => u32::MAX, 1891 (BaseRelayFilterLimitMode::Override(limit), _) => limit, 1892 (_, limit) => limit, 1893 }; 1894 let ids = filter.ids().collect::<Vec<_>>(); 1895 let authors = filter.authors().collect::<Vec<_>>(); 1896 let kinds = filter.kinds().collect::<Vec<_>>(); 1897 let since = 1898 (filter.since() != tangle_store_pocket::PocketTime::min()).then(|| filter.since()); 1899 let until = 1900 (filter.until() != tangle_store_pocket::PocketTime::max()).then(|| filter.until()); 1901 let limit = (limit != u32::MAX).then_some(limit); 1902 PocketOwnedFilter::new( 1903 &ids, 1904 &authors, 1905 &kinds, 1906 filter 1907 .tags() 1908 .map_err(|error| BaseRelayError::error(error.to_string()))?, 1909 since, 1910 until, 1911 limit, 1912 ) 1913 .map_err(|error| BaseRelayError::error(error.to_string())) 1914 } 1915 1916 pub(crate) fn sort_and_dedupe_query_events( 1917 mut events: Vec<PocketOwnedEvent>, 1918 ) -> Vec<PocketOwnedEvent> { 1919 events.sort_by(|left, right| { 1920 let left: &PocketEvent = left; 1921 let right: &PocketEvent = right; 1922 right 1923 .created_at() 1924 .cmp(&left.created_at()) 1925 .then_with(|| left.id().cmp(&right.id())) 1926 }); 1927 let mut seen = BTreeSet::new(); 1928 events 1929 .into_iter() 1930 .filter(|event| { 1931 let event: &PocketEvent = event; 1932 seen.insert(event.id()) 1933 }) 1934 .collect() 1935 } 1936 1937 pub(crate) fn group_read_gate_visible_to_auth( 1938 groups: Option<&GroupServiceHandle>, 1939 event: &(impl GroupEventView + ?Sized), 1940 auth: &GroupAuthContext, 1941 ) -> Result<bool, BaseRelayError> { 1942 groups 1943 .map(|groups| groups.event_visible_to_auth(event, auth)) 1944 .unwrap_or(Ok(true)) 1945 .map_err(BaseRelayError::from) 1946 } 1947 } 1948 1949 pub(crate) fn matched_filter_context( 1950 filter_index: usize, 1951 filter: &PocketFilter, 1952 ) -> BaseRelayMatchedFilterContext { 1953 BaseRelayMatchedFilterContext::from_filter(filter_index, filter) 1954 } 1955 1956 fn pocket_filters_are_complete(filters: &[PocketOwnedFilter]) -> bool { 1957 !filters.is_empty() && filters.iter().all(|filter| filter.completes()) 1958 } 1959 1960 #[cfg(test)] 1961 mod tests { 1962 use super::{ 1963 BaseRelay, BaseRelayCountHll, BaseRelayCountHllTargetPolicy, BaseRelayLimitSettings, 1964 BaseRelayLimits, NEGENTROPY_DISABLED_MESSAGE, 1965 }; 1966 use crate::pocket_conversion::{tangle_event_to_pocket, tangle_filter_to_pocket}; 1967 use crate::relay::auth::BaseAuthState; 1968 use crate::relay::live::CloseResult; 1969 use tangle_crypto::RelaySigner; 1970 use tangle_groups::{ 1971 GroupAuthContext, GroupId, KIND_GROUP_ADMINS, KIND_GROUP_CREATE_GROUP, 1972 KIND_GROUP_CREATE_INVITE, KIND_GROUP_DELETE_EVENT, KIND_GROUP_DELETE_GROUP, 1973 KIND_GROUP_EDIT_METADATA, KIND_GROUP_JOIN_REQUEST, KIND_GROUP_LEAVE_REQUEST, 1974 KIND_GROUP_MEMBERS, KIND_GROUP_METADATA, KIND_GROUP_PUT_USER, KIND_GROUP_REMOVE_USER, 1975 MemberStatus, NIP29_RELAY_GENERATED_KIND_VALUES, StoreOffset, 1976 parse_group_runtime_config_json, 1977 }; 1978 use tangle_protocol::{ 1979 ClientMessage, Event, EventId, Filter, Kind, PublicKeyHex, RelayMessage, SignatureHex, 1980 SubscriptionId, Tag, UnixTimestamp, UnsignedEvent, filter_from_value, 1981 }; 1982 use tangle_store_pocket::{ 1983 PocketEvent, PocketHll8, PocketKind, PocketOwnedEvent, PocketOwnedFilter, PocketOwnedTags, 1984 PocketQueryConfig, PocketStoreConfig, PocketSyncPolicy, PocketTime, 1985 }; 1986 1987 trait BaseRelayCountTestExt { 1988 fn handle_count_protocol( 1989 &self, 1990 subscription_id: SubscriptionId, 1991 filters: Vec<Filter>, 1992 ) -> Result<RelayMessage, crate::errors::BaseRelayError>; 1993 1994 fn handle_count_with_auth_protocol( 1995 &self, 1996 subscription_id: SubscriptionId, 1997 filters: Vec<Filter>, 1998 auth: &BaseAuthState, 1999 ) -> Result<RelayMessage, crate::errors::BaseRelayError>; 2000 } 2001 2002 impl BaseRelayCountTestExt for BaseRelay { 2003 fn handle_count_protocol( 2004 &self, 2005 subscription_id: SubscriptionId, 2006 filters: Vec<Filter>, 2007 ) -> Result<RelayMessage, crate::errors::BaseRelayError> { 2008 let search_present = filters.iter().any(|filter| filter.search().is_some()); 2009 let filters = pocket_filters(filters)?; 2010 self.handle_count_with_group_auth_report( 2011 subscription_id, 2012 filters, 2013 search_present, 2014 &GroupAuthContext::unauthenticated(), 2015 ) 2016 .map(|report| report.into_message()) 2017 } 2018 2019 fn handle_count_with_auth_protocol( 2020 &self, 2021 subscription_id: SubscriptionId, 2022 filters: Vec<Filter>, 2023 auth: &BaseAuthState, 2024 ) -> Result<RelayMessage, crate::errors::BaseRelayError> { 2025 let search_present = filters.iter().any(|filter| filter.search().is_some()); 2026 let filters = pocket_filters(filters)?; 2027 let group_auth = GroupAuthContext::new(auth.authenticated_pubkeys().iter().cloned()); 2028 self.handle_count_with_group_auth_report( 2029 subscription_id, 2030 filters, 2031 search_present, 2032 &group_auth, 2033 ) 2034 .map(|report| report.into_message()) 2035 } 2036 } 2037 2038 fn pocket_filters( 2039 filters: Vec<Filter>, 2040 ) -> Result<Vec<PocketOwnedFilter>, crate::errors::BaseRelayError> { 2041 filters.iter().map(tangle_filter_to_pocket).collect() 2042 } 2043 2044 #[test] 2045 fn base_relay_stores_queries_counts_closes_and_fans_out_public_events() { 2046 let mut relay = test_relay("base-relay-public", 4); 2047 let event = signed_public_event(7, 1, Vec::new(), "hello"); 2048 let subscription_id = SubscriptionId::new("sub-a").expect("sub"); 2049 let filter = filter_from_value(&serde_json::json!({"kinds":[1]})).expect("filter"); 2050 2051 assert_eq!( 2052 relay.handle_event(event.clone()).expect("event"), 2053 RelayMessage::Ok { 2054 event_id: event.id().clone(), 2055 accepted: true, 2056 message: String::new() 2057 } 2058 ); 2059 assert_eq!( 2060 relay.handle_event(event.clone()).expect("duplicate"), 2061 RelayMessage::Ok { 2062 event_id: event.id().clone(), 2063 accepted: true, 2064 message: "duplicate: already have this event".to_owned() 2065 } 2066 ); 2067 2068 let messages = relay 2069 .handle_protocol_req_for_test(subscription_id.clone(), vec![filter.clone()]) 2070 .expect("req"); 2071 assert!( 2072 matches!(&messages[0], RelayMessage::Event { event: found, .. } if found.id() == event.id()) 2073 ); 2074 assert_eq!(messages[1], RelayMessage::Eose(subscription_id.clone())); 2075 assert_eq!( 2076 relay 2077 .handle_count_protocol(subscription_id.clone(), vec![filter]) 2078 .expect("count"), 2079 RelayMessage::Count { 2080 subscription_id: subscription_id.clone(), 2081 count: 1, 2082 hll: None 2083 } 2084 ); 2085 assert!(matches!( 2086 relay.fanout_protocol_for_test(&event).as_slice(), 2087 [RelayMessage::Event { subscription_id: delivered, event: found }] 2088 if delivered == &subscription_id && found.id() == event.id() 2089 )); 2090 assert_eq!(relay.handle_close(&subscription_id), CloseResult::Closed); 2091 assert_eq!(relay.active_subscription_count(), 0); 2092 assert!(relay.fanout_protocol_for_test(&event).is_empty()); 2093 } 2094 2095 #[test] 2096 fn base_relay_uses_configured_pocket_query_scrape_controls() { 2097 let strict_config = test_store_config("base-relay-query-strict"); 2098 let mut strict = BaseRelay::open( 2099 &strict_config, 2100 relay_limits(4), 2101 PocketQueryConfig::new(false, 0, 0), 2102 ) 2103 .expect("strict"); 2104 let strict_event = signed_public_event(7, 1, Vec::new(), "strict"); 2105 let broad = filter_from_value(&serde_json::json!({"limit":1})).expect("filter"); 2106 2107 assert_accepted( 2108 strict 2109 .handle_event(strict_event.clone()) 2110 .expect("strict event"), 2111 &strict_event, 2112 ); 2113 assert!( 2114 strict 2115 .handle_protocol_req_for_test( 2116 SubscriptionId::new("strict").expect("sub"), 2117 vec![broad.clone()] 2118 ) 2119 .expect_err("strict scrape") 2120 .prefixed_message() 2121 .to_lowercase() 2122 .contains("scraper") 2123 ); 2124 2125 let limited_config = test_store_config("base-relay-query-limited"); 2126 let mut limited = BaseRelay::open( 2127 &limited_config, 2128 relay_limits(4), 2129 PocketQueryConfig::new(false, 1, 0), 2130 ) 2131 .expect("limited"); 2132 let limited_event = signed_public_event(8, 1, Vec::new(), "limited"); 2133 2134 assert_accepted( 2135 limited 2136 .handle_event(limited_event.clone()) 2137 .expect("limited event"), 2138 &limited_event, 2139 ); 2140 let messages = limited 2141 .handle_protocol_req_for_test(SubscriptionId::new("limited").expect("sub"), vec![broad]) 2142 .expect("limited scrape"); 2143 2144 assert!( 2145 matches!(&messages[0], RelayMessage::Event { event, .. } if event.id() == limited_event.id()) 2146 ); 2147 } 2148 2149 #[test] 2150 fn base_relay_rejects_search_req_and_count_as_unsupported() { 2151 let mut relay = test_relay("base-relay-search-unsupported", 4); 2152 let req_id = SubscriptionId::new("search-req").expect("req"); 2153 let count_id = SubscriptionId::new("search-count").expect("count"); 2154 let search = filter_from_value(&serde_json::json!({ 2155 "search": "fresh carrots", 2156 "limit": 1 2157 })) 2158 .expect("filter"); 2159 2160 assert_eq!( 2161 relay 2162 .handle_protocol_req_for_test(req_id.clone(), vec![search.clone()]) 2163 .expect("req"), 2164 vec![RelayMessage::Closed { 2165 subscription_id: req_id, 2166 message: "unsupported: search filters are not supported".to_owned() 2167 }] 2168 ); 2169 assert_eq!(relay.active_subscription_count(), 0); 2170 assert_eq!( 2171 relay 2172 .handle_count_protocol(count_id.clone(), vec![search]) 2173 .expect("count"), 2174 RelayMessage::Closed { 2175 subscription_id: count_id, 2176 message: "unsupported: search filters are not supported".to_owned() 2177 } 2178 ); 2179 } 2180 2181 #[test] 2182 fn base_relay_dispatch_returns_disabled_negentropy_surface() { 2183 let mut relay = test_relay("base-relay-negentropy-disabled", 4); 2184 let mut auth = 2185 BaseAuthState::new("wss://relay.radroots.test", 60, 600).expect("auth state"); 2186 let subscription_id = SubscriptionId::new("neg-sub").expect("sub"); 2187 2188 assert_eq!( 2189 relay 2190 .handle_client_message( 2191 ClientMessage::NegOpen { 2192 subscription_id: subscription_id.clone(), 2193 filter: Filter::empty(), 2194 message: "00".to_owned() 2195 }, 2196 &mut auth, 2197 UnixTimestamp::new(100) 2198 ) 2199 .expect("neg open"), 2200 vec![RelayMessage::NegErr { 2201 subscription_id: subscription_id.clone(), 2202 message: NEGENTROPY_DISABLED_MESSAGE.to_owned() 2203 }] 2204 ); 2205 assert_eq!( 2206 relay 2207 .handle_client_message( 2208 ClientMessage::NegMsg { 2209 subscription_id: subscription_id.clone(), 2210 message: String::new() 2211 }, 2212 &mut auth, 2213 UnixTimestamp::new(101) 2214 ) 2215 .expect("neg msg"), 2216 vec![RelayMessage::NegErr { 2217 subscription_id: subscription_id.clone(), 2218 message: NEGENTROPY_DISABLED_MESSAGE.to_owned() 2219 }] 2220 ); 2221 assert_eq!( 2222 relay 2223 .handle_client_message( 2224 ClientMessage::NegClose(subscription_id), 2225 &mut auth, 2226 UnixTimestamp::new(102) 2227 ) 2228 .expect("neg close"), 2229 Vec::<RelayMessage>::new() 2230 ); 2231 } 2232 2233 #[test] 2234 fn base_relay_disabled_negentropy_does_not_validate_or_screen_filter() { 2235 let owner = signer(7).public_key().clone(); 2236 let owner_auth = authenticated_state(7); 2237 let mut auth = 2238 BaseAuthState::new("wss://relay.radroots.test", 60, 600).expect("auth state"); 2239 let mut relay = test_relay_with_groups( 2240 "base-relay-negentropy-disabled-no-screen", 2241 4, 2242 &enabled_groups_for_owner(&owner), 2243 ); 2244 let private_create = signed_private_group_create_event(7, "PrivateNegentropy"); 2245 assert_accepted( 2246 relay 2247 .handle_event_with_auth(private_create.clone(), &owner_auth) 2248 .expect("private create"), 2249 &private_create, 2250 ); 2251 let private_event = signed_event_at( 2252 7, 2253 1, 2254 vec![h("PrivateNegentropy")], 2255 "private negentropy", 2256 1_714_124_434, 2257 ); 2258 assert_accepted( 2259 relay 2260 .handle_event_with_auth(private_event.clone(), &owner_auth) 2261 .expect("private event"), 2262 &private_event, 2263 ); 2264 let subscription_id = SubscriptionId::new("neg-noscreen").expect("sub"); 2265 let filter = filter_from_value(&serde_json::json!({ 2266 "kinds": [1], 2267 "#h": ["PrivateNegentropy"], 2268 "limit": 501 2269 })) 2270 .expect("filter"); 2271 2272 assert_eq!( 2273 relay 2274 .handle_client_message( 2275 ClientMessage::NegOpen { 2276 subscription_id: subscription_id.clone(), 2277 filter, 2278 message: "00".to_owned() 2279 }, 2280 &mut auth, 2281 UnixTimestamp::new(100) 2282 ) 2283 .expect("neg open"), 2284 vec![RelayMessage::NegErr { 2285 subscription_id: subscription_id.clone(), 2286 message: NEGENTROPY_DISABLED_MESSAGE.to_owned() 2287 }] 2288 ); 2289 assert_eq!( 2290 relay 2291 .handle_client_message( 2292 ClientMessage::NegMsg { 2293 subscription_id: subscription_id.clone(), 2294 message: "should-not-touch-storage".to_owned() 2295 }, 2296 &mut auth, 2297 UnixTimestamp::new(101) 2298 ) 2299 .expect("neg msg"), 2300 vec![RelayMessage::NegErr { 2301 subscription_id, 2302 message: NEGENTROPY_DISABLED_MESSAGE.to_owned() 2303 }] 2304 ); 2305 } 2306 2307 #[test] 2308 fn base_relay_fetches_events_by_store_offset() { 2309 let relay = test_relay("base-relay-offset-lookup", 4); 2310 let event = signed_public_event(7, 1, Vec::new(), "offset"); 2311 let pocket = tangle_event_to_pocket(&event).expect("pocket"); 2312 let offset = StoreOffset::new(relay.store.store_event(&pocket).expect("store")); 2313 2314 let found = relay.event_by_offset(offset).expect("offset"); 2315 let found: &PocketEvent = &found; 2316 assert_eq!(found.id().as_hex_string(), event.id().as_str()); 2317 } 2318 2319 #[test] 2320 fn base_relay_req_merges_filters_with_order_dedupe_and_limits() { 2321 let mut relay = test_relay("base-relay-req-order", 8); 2322 let market_tag = Tag::from_parts("t", &["market"]).expect("tag"); 2323 let old_market = 2324 signed_event_at(7, 1, vec![market_tag.clone()], "old market", 1_714_124_433); 2325 let tied_author = 2326 signed_event_at(7, 1, vec![market_tag.clone()], "tied author", 1_714_124_434); 2327 let tied_other = 2328 signed_event_at(8, 1, vec![market_tag.clone()], "tied other", 1_714_124_434); 2329 let kind_two = signed_event_at(7, 2, Vec::new(), "kind two", 1_714_124_435); 2330 let wrong_tag = signed_event_at( 2331 9, 2332 1, 2333 vec![Tag::from_parts("t", &["other"]).expect("tag")], 2334 "wrong tag", 2335 1_714_124_436, 2336 ); 2337 2338 for event in [ 2339 &old_market, 2340 &tied_other, 2341 &kind_two, 2342 &wrong_tag, 2343 &tied_author, 2344 ] { 2345 assert_accepted(relay.handle_event(event.clone()).expect("event"), event); 2346 } 2347 2348 let subscription_id = SubscriptionId::new("req-order").expect("sub"); 2349 let market_limit = 2350 filter_from_value(&serde_json::json!({"kinds":[1],"#t":["market"],"limit":2})) 2351 .expect("market filter"); 2352 let author_limit = filter_from_value(&serde_json::json!({ 2353 "authors":[tied_author.unsigned().pubkey().as_str()], 2354 "kinds":[1,2], 2355 "limit":2 2356 })) 2357 .expect("author filter"); 2358 let messages = relay 2359 .handle_protocol_req_for_test(subscription_id.clone(), vec![market_limit, author_limit]) 2360 .expect("req"); 2361 let mut tied = [tied_author.clone(), tied_other.clone()]; 2362 tied.sort_by(|left, right| left.id().cmp(right.id())); 2363 let expected = [kind_two.clone(), tied[0].clone(), tied[1].clone()]; 2364 2365 assert_eq!(messages.len(), expected.len() + 1); 2366 for (message, event) in messages.iter().zip(expected.iter()) { 2367 assert!(matches!( 2368 message, 2369 RelayMessage::Event { 2370 subscription_id: actual, 2371 event: found 2372 } if actual == &subscription_id && found.id() == event.id() 2373 )); 2374 } 2375 assert_eq!( 2376 messages.last(), 2377 Some(&RelayMessage::Eose(subscription_id.clone())) 2378 ); 2379 assert!(!messages.iter().any(|message| matches!( 2380 message, 2381 RelayMessage::Event { event, .. } 2382 if event.id() == old_market.id() || event.id() == wrong_tag.id() 2383 ))); 2384 } 2385 2386 #[test] 2387 fn base_relay_req_count_paths_preserve_chorus_parity() { 2388 let owner = signer(7).public_key().clone(); 2389 let auth = authenticated_state(7); 2390 let outsider_auth = authenticated_state(8); 2391 let mut relay = test_relay_with_groups( 2392 "base-relay-req-count-chorus-parity", 2393 8, 2394 &enabled_groups_for_owner(&owner), 2395 ); 2396 let market_tag = Tag::from_parts("t", &["market"]).expect("tag"); 2397 let old_market = 2398 signed_event_at(7, 1, vec![market_tag.clone()], "old market", 1_714_124_433); 2399 let tied_author = 2400 signed_event_at(7, 1, vec![market_tag.clone()], "tied author", 1_714_124_434); 2401 let tied_other = 2402 signed_event_at(8, 1, vec![market_tag.clone()], "tied other", 1_714_124_434); 2403 let kind_two = signed_event_at(7, 2, Vec::new(), "kind two", 1_714_124_435); 2404 let wrong_tag = signed_event_at( 2405 9, 2406 1, 2407 vec![Tag::from_parts("t", &["other"]).expect("tag")], 2408 "wrong tag", 2409 1_714_124_436, 2410 ); 2411 for event in [ 2412 &old_market, 2413 &tied_other, 2414 &kind_two, 2415 &wrong_tag, 2416 &tied_author, 2417 ] { 2418 assert_accepted(relay.handle_event(event.clone()).expect("event"), event); 2419 } 2420 relay 2421 .handle_event_with_auth(signed_private_group_create_event(7, "Private"), &auth) 2422 .expect("private create"); 2423 let private_market = signed_event_at( 2424 7, 2425 1, 2426 vec![h("Private"), market_tag.clone()], 2427 "private market", 2428 1_714_124_437, 2429 ); 2430 assert_accepted( 2431 relay 2432 .handle_event_with_auth(private_market.clone(), &auth) 2433 .expect("private event"), 2434 &private_market, 2435 ); 2436 2437 let subscription_id = SubscriptionId::new("req-count-parity").expect("sub"); 2438 let market_limit = 2439 filter_from_value(&serde_json::json!({"kinds":[1],"#t":["market"],"limit":2})) 2440 .expect("market filter"); 2441 let author_limit = filter_from_value(&serde_json::json!({ 2442 "authors":[tied_author.unsigned().pubkey().as_str()], 2443 "kinds":[1,2], 2444 "limit":2 2445 })) 2446 .expect("author filter"); 2447 let messages = relay 2448 .handle_protocol_req_for_test( 2449 subscription_id.clone(), 2450 vec![market_limit.clone(), author_limit.clone()], 2451 ) 2452 .expect("req"); 2453 let mut tied = [tied_author.clone(), tied_other.clone()]; 2454 tied.sort_by(|left, right| left.id().cmp(right.id())); 2455 let expected = [kind_two.clone(), tied[0].clone(), tied[1].clone()]; 2456 let event_ids = messages 2457 .iter() 2458 .filter_map(|message| match message { 2459 RelayMessage::Event { 2460 subscription_id: actual, 2461 event, 2462 } if actual == &subscription_id => Some(event.id().clone()), 2463 _ => None, 2464 }) 2465 .collect::<Vec<_>>(); 2466 let expected_ids = expected 2467 .iter() 2468 .map(|event| event.id().clone()) 2469 .collect::<Vec<_>>(); 2470 2471 assert_eq!(event_ids, expected_ids); 2472 assert_eq!( 2473 messages.last(), 2474 Some(&RelayMessage::Closed { 2475 subscription_id: subscription_id.clone(), 2476 message: "auth-required: authentication required to read group events".to_owned() 2477 }) 2478 ); 2479 assert!(!messages.iter().any( 2480 |message| matches!(message, RelayMessage::Eose(actual) if actual == &subscription_id) 2481 )); 2482 assert!(!event_ids.contains(private_market.id())); 2483 assert!(!event_ids.contains(old_market.id())); 2484 assert!(!event_ids.contains(wrong_tag.id())); 2485 assert_eq!(relay.active_subscription_count(), 0); 2486 2487 let restricted_sub = SubscriptionId::new("restricted-screened").expect("sub"); 2488 let restricted_messages = relay 2489 .handle_protocol_req_with_auth_for_test( 2490 restricted_sub.clone(), 2491 vec![market_limit.clone(), author_limit.clone()], 2492 &outsider_auth, 2493 ) 2494 .expect("restricted req"); 2495 let restricted_event_ids = restricted_messages 2496 .iter() 2497 .filter_map(|message| match message { 2498 RelayMessage::Event { 2499 subscription_id: actual, 2500 event, 2501 } if actual == &restricted_sub => Some(event.id().clone()), 2502 _ => None, 2503 }) 2504 .collect::<Vec<_>>(); 2505 assert_eq!(restricted_event_ids, expected_ids); 2506 assert_eq!( 2507 restricted_messages.last(), 2508 Some(&RelayMessage::Closed { 2509 subscription_id: restricted_sub.clone(), 2510 message: "restricted: group is unavailable".to_owned() 2511 }) 2512 ); 2513 assert!(!restricted_messages.iter().any( 2514 |message| matches!(message, RelayMessage::Eose(actual) if actual == &restricted_sub) 2515 )); 2516 assert_eq!(relay.active_subscription_count(), 0); 2517 2518 let private_sub = SubscriptionId::new("private-screened").expect("sub"); 2519 assert_eq!( 2520 relay 2521 .handle_protocol_req_for_test( 2522 private_sub.clone(), 2523 vec![filter_group_tag(1, "h", "Private")] 2524 ) 2525 .expect("private unauth req"), 2526 vec![RelayMessage::Closed { 2527 subscription_id: private_sub, 2528 message: "auth-required: authentication required to read group events".to_owned() 2529 }] 2530 ); 2531 assert_eq!(relay.active_subscription_count(), 0); 2532 let private_auth_sub = SubscriptionId::new("private-auth").expect("sub"); 2533 assert!(matches!( 2534 relay 2535 .handle_protocol_req_with_auth_for_test( 2536 private_auth_sub.clone(), 2537 vec![filter_group_tag(1, "h", "Private")], 2538 &auth 2539 ) 2540 .expect("private auth req") 2541 .as_slice(), 2542 [RelayMessage::Event { subscription_id, event }, RelayMessage::Eose(eose)] 2543 if subscription_id == &private_auth_sub && event.id() == private_market.id() && eose == &private_auth_sub 2544 )); 2545 2546 let market_notes = 2547 filter_from_value(&serde_json::json!({"kinds":[1],"#t":["market"],"limit":10})) 2548 .expect("market count filter"); 2549 let author_events = filter_from_value(&serde_json::json!({ 2550 "authors":[tied_author.unsigned().pubkey().as_str()], 2551 "kinds":[1,2], 2552 "limit":10 2553 })) 2554 .expect("author count filter"); 2555 assert_eq!( 2556 relay 2557 .handle_count_protocol( 2558 SubscriptionId::new("count-visible").expect("sub"), 2559 vec![market_notes.clone(), author_events.clone()] 2560 ) 2561 .expect("visible count"), 2562 RelayMessage::Count { 2563 subscription_id: SubscriptionId::new("count-visible").expect("sub"), 2564 count: 4, 2565 hll: None 2566 } 2567 ); 2568 assert_eq!( 2569 relay 2570 .handle_count_with_auth_protocol( 2571 SubscriptionId::new("count-auth").expect("sub"), 2572 vec![market_notes, author_events], 2573 &auth 2574 ) 2575 .expect("auth count"), 2576 RelayMessage::Count { 2577 subscription_id: SubscriptionId::new("count-auth").expect("sub"), 2578 count: 5, 2579 hll: None 2580 } 2581 ); 2582 2583 let too_large_limit = 2584 filter_from_value(&serde_json::json!({"limit":501})).expect("limit filter"); 2585 assert!( 2586 relay 2587 .handle_protocol_req_for_test( 2588 SubscriptionId::new("limit-req").expect("sub"), 2589 vec![too_large_limit.clone()] 2590 ) 2591 .expect_err("req limit") 2592 .prefixed_message() 2593 .contains("max_limit 500") 2594 ); 2595 assert!( 2596 relay 2597 .handle_count_protocol( 2598 SubscriptionId::new("limit-count").expect("sub"), 2599 vec![too_large_limit] 2600 ) 2601 .expect_err("count limit") 2602 .prefixed_message() 2603 .contains("max_limit 500") 2604 ); 2605 2606 let search = filter_from_value(&serde_json::json!({"search":"carrots","limit":1})) 2607 .expect("search filter"); 2608 let search_req = SubscriptionId::new("search-req").expect("sub"); 2609 assert_eq!( 2610 relay 2611 .handle_protocol_req_for_test(search_req.clone(), vec![search.clone()]) 2612 .expect("search req"), 2613 vec![RelayMessage::Closed { 2614 subscription_id: search_req, 2615 message: "unsupported: search filters are not supported".to_owned() 2616 }] 2617 ); 2618 let search_count = SubscriptionId::new("search-count").expect("sub"); 2619 assert_eq!( 2620 relay 2621 .handle_count_protocol(search_count.clone(), vec![search]) 2622 .expect("search count"), 2623 RelayMessage::Closed { 2624 subscription_id: search_count, 2625 message: "unsupported: search filters are not supported".to_owned() 2626 } 2627 ); 2628 } 2629 2630 #[test] 2631 fn base_relay_enforces_runtime_limits() { 2632 let config = test_store_config("base-relay-runtime-limits"); 2633 let mut relay = BaseRelay::open( 2634 &config, 2635 BaseRelayLimits::new(BaseRelayLimitSettings { 2636 max_pending_events: 2, 2637 max_subscription_id_length: 3, 2638 max_subscriptions: 1, 2639 max_filters_per_request: 1, 2640 max_tag_values_per_filter: 1, 2641 max_query_complexity: 4, 2642 max_event_tags: 1, 2643 max_content_length: 4, 2644 max_limit: 2, 2645 default_limit: 1, 2646 }) 2647 .expect("limits"), 2648 PocketQueryConfig::default(), 2649 ) 2650 .expect("relay"); 2651 let first = signed_event_at(7, 1, Vec::new(), "one", 1_714_124_430); 2652 let second = signed_event_at(8, 1, Vec::new(), "two", 1_714_124_431); 2653 2654 assert_accepted(relay.handle_event(first.clone()).expect("first"), &first); 2655 assert_accepted(relay.handle_event(second.clone()).expect("second"), &second); 2656 2657 let limited = relay 2658 .handle_protocol_req_for_test( 2659 SubscriptionId::new("lim").expect("sub"), 2660 vec![Filter::empty()], 2661 ) 2662 .expect("limited"); 2663 assert_eq!( 2664 limited 2665 .iter() 2666 .filter(|message| matches!(message, RelayMessage::Event { .. })) 2667 .count(), 2668 1 2669 ); 2670 assert_eq!( 2671 relay.handle_close(&SubscriptionId::new("lim").expect("sub")), 2672 CloseResult::Closed 2673 ); 2674 2675 assert!( 2676 relay 2677 .handle_protocol_req_for_test( 2678 SubscriptionId::new("long").expect("sub"), 2679 vec![Filter::empty()] 2680 ) 2681 .expect_err("subscription id length") 2682 .prefixed_message() 2683 .contains("max_subid_length 3") 2684 ); 2685 assert!( 2686 relay 2687 .handle_count_protocol( 2688 SubscriptionId::new("cnt").expect("sub"), 2689 vec![Filter::empty(), Filter::empty()] 2690 ) 2691 .expect_err("filter count") 2692 .prefixed_message() 2693 .contains("max_filters_per_request 1") 2694 ); 2695 assert!( 2696 relay 2697 .handle_count_protocol( 2698 SubscriptionId::new("tag").expect("sub"), 2699 vec![ 2700 filter_from_value(&serde_json::json!({"#t":["one", "two"]})) 2701 .expect("filter") 2702 ] 2703 ) 2704 .expect_err("tag values") 2705 .prefixed_message() 2706 .contains("max_tag_values_per_filter 1") 2707 ); 2708 assert!( 2709 relay 2710 .handle_count_protocol( 2711 SubscriptionId::new("max").expect("sub"), 2712 vec![filter_from_value(&serde_json::json!({"limit":3})).expect("filter")] 2713 ) 2714 .expect_err("max limit") 2715 .prefixed_message() 2716 .contains("max_limit 2") 2717 ); 2718 2719 let too_many_tags = signed_event_at( 2720 9, 2721 1, 2722 vec![ 2723 Tag::from_parts("t", &["one"]).expect("tag"), 2724 Tag::from_parts("p", &["two"]).expect("tag"), 2725 ], 2726 "ok", 2727 1_714_124_432, 2728 ); 2729 assert!(matches!( 2730 relay.handle_event(too_many_tags).expect("tags"), 2731 RelayMessage::Ok { accepted: false, message, .. } 2732 if message.contains("max_event_tags 1") 2733 )); 2734 2735 let too_much_content = signed_event_at(10, 1, Vec::new(), "12345", 1_714_124_433); 2736 assert!(matches!( 2737 relay.handle_event(too_much_content).expect("content"), 2738 RelayMessage::Ok { accepted: false, message, .. } 2739 if message.contains("max_content_length 4") 2740 )); 2741 } 2742 2743 #[test] 2744 fn base_relay_rejects_over_budget_req_and_count() { 2745 let config = test_store_config("base-relay-query-complexity"); 2746 let mut relay = BaseRelay::open( 2747 &config, 2748 BaseRelayLimits::new(BaseRelayLimitSettings { 2749 max_pending_events: 4, 2750 max_subscription_id_length: 64, 2751 max_subscriptions: 64, 2752 max_filters_per_request: 10, 2753 max_tag_values_per_filter: 10, 2754 max_query_complexity: 4, 2755 max_event_tags: 200, 2756 max_content_length: 65_536, 2757 max_limit: 10, 2758 default_limit: 1, 2759 }) 2760 .expect("limits"), 2761 PocketQueryConfig::default(), 2762 ) 2763 .expect("relay"); 2764 let complex = filter_from_value(&serde_json::json!({ 2765 "kinds": [1], 2766 "#t": ["market"], 2767 "limit": 2 2768 })) 2769 .expect("filter"); 2770 2771 assert!( 2772 relay 2773 .handle_protocol_req_for_test( 2774 SubscriptionId::new("req").expect("sub"), 2775 vec![complex.clone()] 2776 ) 2777 .expect_err("req complexity") 2778 .prefixed_message() 2779 .contains("max_query_complexity 4") 2780 ); 2781 assert_eq!(relay.active_subscription_count(), 0); 2782 assert!( 2783 relay 2784 .handle_count_protocol(SubscriptionId::new("cnt").expect("sub"), vec![complex]) 2785 .expect_err("count complexity") 2786 .prefixed_message() 2787 .contains("max_query_complexity 4") 2788 ); 2789 } 2790 2791 #[test] 2792 fn base_relay_count_dedupes_overlapping_visible_filters() { 2793 let relay = test_relay("base-relay-count-dedupe", 8); 2794 let market_tag = Tag::from_parts("t", &["market"]).expect("tag"); 2795 let first = signed_event_at(7, 1, vec![market_tag.clone()], "first", 1_714_124_433); 2796 let second = signed_event_at(8, 1, vec![market_tag], "second", 1_714_124_434); 2797 let third = signed_event_at(7, 2, Vec::new(), "third", 1_714_124_435); 2798 2799 for event in [&first, &second, &third] { 2800 assert_accepted(relay.handle_event(event.clone()).expect("event"), event); 2801 } 2802 2803 let market_notes = 2804 filter_from_value(&serde_json::json!({"kinds":[1],"#t":["market"],"limit":2})) 2805 .expect("market filter"); 2806 let author_events = filter_from_value(&serde_json::json!({ 2807 "authors":[first.unsigned().pubkey().as_str()], 2808 "kinds":[1,2], 2809 "limit":10 2810 })) 2811 .expect("author filter"); 2812 let limited_market = 2813 filter_from_value(&serde_json::json!({"kinds":[1],"#t":["market"],"limit":1})) 2814 .expect("limited filter"); 2815 2816 assert_eq!( 2817 relay 2818 .handle_count_protocol( 2819 SubscriptionId::new("count-limit").expect("sub"), 2820 vec![limited_market] 2821 ) 2822 .expect("count"), 2823 RelayMessage::Count { 2824 subscription_id: SubscriptionId::new("count-limit").expect("sub"), 2825 count: 2, 2826 hll: None 2827 } 2828 ); 2829 2830 assert_eq!( 2831 relay 2832 .handle_count_protocol( 2833 SubscriptionId::new("count-dedupe").expect("sub"), 2834 vec![market_notes, author_events] 2835 ) 2836 .expect("count"), 2837 RelayMessage::Count { 2838 subscription_id: SubscriptionId::new("count-dedupe").expect("sub"), 2839 count: 3, 2840 hll: None 2841 } 2842 ); 2843 } 2844 2845 #[test] 2846 fn base_relay_count_hll_emits_for_public_single_filter() { 2847 let relay = test_relay("base-relay-count-hll-public", 8); 2848 let target = "a".repeat(EventId::HEX_LENGTH); 2849 let target_tag = Tag::from_parts("e", &[&target]).expect("tag"); 2850 let first = signed_pocket_public_event(7, 7, vec![target_tag.clone()], "first reaction"); 2851 let second = signed_pocket_public_event(8, 7, vec![target_tag], "second reaction"); 2852 2853 for event in [&first, &second] { 2854 assert_pocket_accepted(relay.handle_pocket_event(event).expect("event"), event); 2855 } 2856 2857 let RelayMessage::Count { count, hll, .. } = relay 2858 .handle_count_protocol( 2859 SubscriptionId::new("count-hll-public").expect("sub"), 2860 vec![ 2861 filter_from_value(&serde_json::json!({"kinds":[7],"#e":[target]})) 2862 .expect("filter"), 2863 ], 2864 ) 2865 .expect("count") 2866 else { 2867 panic!("count expected") 2868 }; 2869 let hll = hll.expect("hll"); 2870 2871 assert_eq!(count, 2); 2872 assert_eq!(hll.len(), 512); 2873 assert_ne!(hll, "00".repeat(256)); 2874 } 2875 2876 #[test] 2877 fn base_relay_count_hll_omits_for_private_hidden_unknown_limited_multi_and_redacted_counts() { 2878 let owner = signer(7).public_key().clone(); 2879 let owner_auth = authenticated_state(7); 2880 let unauth = BaseAuthState::new("wss://relay.radroots.test", 60, 600).expect("auth state"); 2881 let relay = test_relay_with_groups( 2882 "base-relay-count-hll-omits", 2883 8, 2884 &enabled_groups_for_owner(&owner), 2885 ); 2886 let target = "b".repeat(EventId::HEX_LENGTH); 2887 let target_tag = Tag::from_parts("e", &[&target]).expect("tag"); 2888 let public = signed_pocket_public_event(8, 7, vec![target_tag.clone()], "public reaction"); 2889 2890 assert_pocket_accepted(relay.handle_pocket_event(&public).expect("public"), &public); 2891 let private_create = signed_pocket_private_group_create_event(7, "PrivateHll"); 2892 assert_pocket_accepted( 2893 relay 2894 .handle_pocket_event_with_auth(&private_create, &owner_auth) 2895 .expect("private create"), 2896 &private_create, 2897 ); 2898 let private = signed_pocket_event_at_tags( 2899 7, 2900 7, 2901 vec![h("PrivateHll"), target_tag.clone()], 2902 "private reaction", 2903 1_714_124_434, 2904 ); 2905 assert_pocket_accepted( 2906 relay 2907 .handle_pocket_event_with_auth(&private, &owner_auth) 2908 .expect("private reaction"), 2909 &private, 2910 ); 2911 let hidden_create = signed_pocket_group_create_event_with_tags( 2912 7, 2913 "HiddenHll", 2914 vec![hidden()], 2915 1_714_124_435, 2916 ); 2917 assert_pocket_accepted( 2918 relay 2919 .handle_pocket_event_with_auth(&hidden_create, &owner_auth) 2920 .expect("hidden create"), 2921 &hidden_create, 2922 ); 2923 let hidden = signed_pocket_event_at_tags( 2924 7, 2925 7, 2926 vec![h("HiddenHll"), target_tag.clone()], 2927 "hidden reaction", 2928 1_714_124_436, 2929 ); 2930 assert_pocket_accepted( 2931 relay 2932 .handle_pocket_event_with_auth(&hidden, &owner_auth) 2933 .expect("hidden reaction"), 2934 &hidden, 2935 ); 2936 let deleted_create = signed_pocket_group_create_event(7, "DeletedHll"); 2937 assert_pocket_accepted( 2938 relay 2939 .handle_pocket_event_with_auth(&deleted_create, &owner_auth) 2940 .expect("deleted create"), 2941 &deleted_create, 2942 ); 2943 let deleted = signed_pocket_event_at_tags( 2944 7, 2945 7, 2946 vec![h("DeletedHll"), target_tag.clone()], 2947 "deleted reaction", 2948 1_714_124_438, 2949 ); 2950 assert_pocket_accepted( 2951 relay 2952 .handle_pocket_event_with_auth(&deleted, &owner_auth) 2953 .expect("deleted reaction"), 2954 &deleted, 2955 ); 2956 let delete_group = signed_pocket_event_at_tags( 2957 7, 2958 KIND_GROUP_DELETE_GROUP, 2959 vec![h("DeletedHll")], 2960 "", 2961 1_714_124_439, 2962 ); 2963 assert_pocket_accepted( 2964 relay 2965 .handle_pocket_event_with_auth(&delete_group, &owner_auth) 2966 .expect("delete group"), 2967 &delete_group, 2968 ); 2969 let unknown = signed_pocket_event_at_tags( 2970 7, 2971 7, 2972 vec![h("UnknownHll"), target_tag.clone()], 2973 "unknown reaction", 2974 1_714_124_440, 2975 ); 2976 relay.store.store_event(&unknown).expect("store unknown"); 2977 2978 let authorized_private = relay 2979 .handle_count_with_auth_protocol( 2980 SubscriptionId::new("count-hll-authorized-private").expect("sub"), 2981 vec![ 2982 filter_from_value(&serde_json::json!({"kinds":[7],"#e":[target.clone()]})) 2983 .expect("filter"), 2984 ], 2985 &owner_auth, 2986 ) 2987 .expect("authorized private count"); 2988 assert!(matches!( 2989 authorized_private, 2990 RelayMessage::Count { 2991 count: 3, 2992 hll: None, 2993 .. 2994 } 2995 )); 2996 2997 let multi_filter = relay 2998 .handle_count_with_auth_protocol( 2999 SubscriptionId::new("count-hll-multi-filter").expect("sub"), 3000 vec![ 3001 filter_from_value(&serde_json::json!({"kinds":[7],"#e":[target.clone()]})) 3002 .expect("filter"), 3003 filter_from_value(&serde_json::json!({"kinds":[7],"#e":["c".repeat(64)]})) 3004 .expect("filter"), 3005 ], 3006 &owner_auth, 3007 ) 3008 .expect("multi count"); 3009 assert!(matches!( 3010 multi_filter, 3011 RelayMessage::Count { 3012 count: 3, 3013 hll: None, 3014 .. 3015 } 3016 )); 3017 3018 let limited = relay 3019 .handle_count_with_auth_protocol( 3020 SubscriptionId::new("count-hll-limited").expect("sub"), 3021 vec![ 3022 filter_from_value( 3023 &serde_json::json!({"kinds":[7],"#e":[target.clone()],"limit":1}), 3024 ) 3025 .expect("filter"), 3026 ], 3027 &owner_auth, 3028 ) 3029 .expect("limited count"); 3030 assert!(matches!( 3031 limited, 3032 RelayMessage::Count { 3033 count: 3, 3034 hll: None, 3035 .. 3036 } 3037 )); 3038 3039 let redacted = relay 3040 .handle_count_with_auth_protocol( 3041 SubscriptionId::new("count-hll-redacted").expect("sub"), 3042 vec![ 3043 filter_from_value(&serde_json::json!({"kinds":[7],"#e":[target]})) 3044 .expect("filter"), 3045 ], 3046 &unauth, 3047 ) 3048 .expect("redacted count"); 3049 assert!(matches!( 3050 redacted, 3051 RelayMessage::Count { 3052 count: 1, 3053 hll: None, 3054 .. 3055 } 3056 )); 3057 3058 assert_count_without_hll( 3059 &relay, 3060 "count-hll-private-h-target", 3061 serde_json::json!({"kinds":[7],"#h":["PrivateHll"]}), 3062 None, 3063 0, 3064 ); 3065 assert_count_without_hll( 3066 &relay, 3067 "count-hll-hidden-h-target", 3068 serde_json::json!({"kinds":[7],"#h":["HiddenHll"]}), 3069 None, 3070 0, 3071 ); 3072 assert_count_without_hll( 3073 &relay, 3074 "count-hll-unknown-h-target", 3075 serde_json::json!({"kinds":[7],"#h":["UnknownHll"]}), 3076 None, 3077 0, 3078 ); 3079 assert_count_without_hll( 3080 &relay, 3081 "count-hll-deleted-h-target", 3082 serde_json::json!({"kinds":[7],"#h":["DeletedHll"]}), 3083 None, 3084 0, 3085 ); 3086 assert_count_without_hll( 3087 &relay, 3088 "count-hll-private-d-target", 3089 serde_json::json!({"kinds":[KIND_GROUP_METADATA],"#d":["PrivateHll"]}), 3090 None, 3091 1, 3092 ); 3093 assert_count_without_hll( 3094 &relay, 3095 "count-hll-hidden-d-target", 3096 serde_json::json!({"kinds":[KIND_GROUP_METADATA],"#d":["HiddenHll"]}), 3097 None, 3098 0, 3099 ); 3100 assert_count_without_hll( 3101 &relay, 3102 "count-hll-unknown-d-target", 3103 serde_json::json!({"kinds":[KIND_GROUP_METADATA],"#d":["UnknownHll"]}), 3104 None, 3105 0, 3106 ); 3107 assert_count_without_hll( 3108 &relay, 3109 "count-hll-deleted-d-target", 3110 serde_json::json!({"kinds":[KIND_GROUP_METADATA],"#d":["DeletedHll"]}), 3111 None, 3112 0, 3113 ); 3114 } 3115 3116 #[test] 3117 fn base_relay_count_hll_group_target_policy_classifies_h_and_d_targets() { 3118 let owner = signer(7).public_key().clone(); 3119 let owner_auth = authenticated_state(7); 3120 let relay = test_relay_with_groups( 3121 "base-relay-count-hll-target-policy", 3122 8, 3123 &enabled_groups_for_owner(&owner), 3124 ); 3125 for event in [ 3126 signed_pocket_group_create_event(7, "PublicHll"), 3127 signed_pocket_group_create_event(7, "SecondHll"), 3128 signed_pocket_private_group_create_event(7, "PrivateHll"), 3129 signed_pocket_group_create_event_with_tags( 3130 7, 3131 "HiddenHll", 3132 vec![hidden()], 3133 1_714_124_435, 3134 ), 3135 signed_pocket_group_create_event(7, "DeletedHll"), 3136 ] { 3137 assert_pocket_accepted( 3138 relay 3139 .handle_pocket_event_with_auth(&event, &owner_auth) 3140 .expect("group create"), 3141 &event, 3142 ); 3143 } 3144 let delete_group = signed_pocket_event_at_tags( 3145 7, 3146 KIND_GROUP_DELETE_GROUP, 3147 vec![h("DeletedHll")], 3148 "", 3149 1_714_124_436, 3150 ); 3151 assert_pocket_accepted( 3152 relay 3153 .handle_pocket_event_with_auth(&delete_group, &owner_auth) 3154 .expect("delete group"), 3155 &delete_group, 3156 ); 3157 3158 assert_eq!( 3159 hll_target_policy(&relay, serde_json::json!({"kinds":[7],"#h":["PublicHll"]})), 3160 BaseRelayCountHllTargetPolicy::Eligible 3161 ); 3162 assert_eq!( 3163 hll_target_policy( 3164 &relay, 3165 serde_json::json!({"kinds":[7],"#h":["PublicHll","SecondHll"]}) 3166 ), 3167 BaseRelayCountHllTargetPolicy::Eligible 3168 ); 3169 assert_eq!( 3170 hll_target_policy( 3171 &relay, 3172 serde_json::json!({"kinds":[7],"#h":["PublicHll","PrivateHll"]}) 3173 ), 3174 BaseRelayCountHllTargetPolicy::Suppress 3175 ); 3176 assert_eq!( 3177 hll_target_policy(&relay, serde_json::json!({"kinds":[7],"#h":["HiddenHll"]})), 3178 BaseRelayCountHllTargetPolicy::Suppress 3179 ); 3180 assert_eq!( 3181 hll_target_policy(&relay, serde_json::json!({"kinds":[7],"#h":["DeletedHll"]})), 3182 BaseRelayCountHllTargetPolicy::Suppress 3183 ); 3184 assert_eq!( 3185 hll_target_policy(&relay, serde_json::json!({"kinds":[7],"#h":["UnknownHll"]})), 3186 BaseRelayCountHllTargetPolicy::Suppress 3187 ); 3188 assert_eq!( 3189 hll_target_policy(&relay, serde_json::json!({"kinds":[7],"#h":[""]})), 3190 BaseRelayCountHllTargetPolicy::Suppress 3191 ); 3192 assert_eq!( 3193 hll_target_policy( 3194 &relay, 3195 serde_json::json!({"kinds":[KIND_GROUP_METADATA],"#d":["PublicHll"]}) 3196 ), 3197 BaseRelayCountHllTargetPolicy::Eligible 3198 ); 3199 assert_eq!( 3200 hll_target_policy( 3201 &relay, 3202 serde_json::json!({"kinds":[KIND_GROUP_METADATA],"#d":["PrivateHll"]}) 3203 ), 3204 BaseRelayCountHllTargetPolicy::Suppress 3205 ); 3206 assert_eq!( 3207 hll_target_policy(&relay, serde_json::json!({"#d":["PublicHll"]})), 3208 BaseRelayCountHllTargetPolicy::Suppress 3209 ); 3210 assert_eq!( 3211 hll_target_policy( 3212 &relay, 3213 serde_json::json!({"kinds":[30023],"#d":["PrivateHll"]}) 3214 ), 3215 BaseRelayCountHllTargetPolicy::Eligible 3216 ); 3217 3218 let mut private_hll = count_hll_for_target_policy_test(); 3219 let private_filter = [pocket_filter_from_value( 3220 serde_json::json!({"kinds":[7],"#h":["PrivateHll"]}), 3221 )]; 3222 private_hll.suppress_for_filter_targets(relay.groups.as_ref(), &private_filter); 3223 assert!(private_hll.into_hex().is_none()); 3224 3225 let mut non_group_hll = count_hll_for_target_policy_test(); 3226 let non_group_filter = [pocket_filter_from_value( 3227 serde_json::json!({"kinds":[30023],"#d":["PrivateHll"]}), 3228 )]; 3229 non_group_hll.suppress_for_filter_targets(relay.groups.as_ref(), &non_group_filter); 3230 assert!(non_group_hll.into_hex().is_some()); 3231 } 3232 3233 #[test] 3234 fn base_relay_count_hll_group_target_policy_suppresses_unresolved_group_targets() { 3235 let relay = test_relay("base-relay-count-hll-target-policy-no-groups", 8); 3236 3237 assert_eq!( 3238 hll_target_policy(&relay, serde_json::json!({"kinds":[7],"#h":["PublicHll"]})), 3239 BaseRelayCountHllTargetPolicy::Suppress 3240 ); 3241 assert_eq!( 3242 hll_target_policy(&relay, serde_json::json!({"#d":["PublicHll"]})), 3243 BaseRelayCountHllTargetPolicy::Suppress 3244 ); 3245 assert_eq!( 3246 hll_target_policy( 3247 &relay, 3248 serde_json::json!({"kinds":[30023],"#d":["PublicHll"]}) 3249 ), 3250 BaseRelayCountHllTargetPolicy::Eligible 3251 ); 3252 } 3253 3254 #[test] 3255 fn base_relay_count_does_not_apply_default_or_client_limits() { 3256 let config = test_store_config("base-relay-count-no-default-limit"); 3257 let relay = BaseRelay::open( 3258 &config, 3259 BaseRelayLimits::new(BaseRelayLimitSettings { 3260 max_pending_events: 4, 3261 max_subscription_id_length: 64, 3262 max_subscriptions: 64, 3263 max_filters_per_request: 10, 3264 max_tag_values_per_filter: 10, 3265 max_query_complexity: 4, 3266 max_event_tags: 200, 3267 max_content_length: 65_536, 3268 max_limit: 10, 3269 default_limit: 1, 3270 }) 3271 .expect("limits"), 3272 PocketQueryConfig::default(), 3273 ) 3274 .expect("relay"); 3275 let first = signed_event_at(7, 1, Vec::new(), "first", 1_714_124_433); 3276 let second = signed_event_at(7, 1, Vec::new(), "second", 1_714_124_434); 3277 let third = signed_event_at(7, 1, Vec::new(), "third", 1_714_124_435); 3278 3279 for event in [&first, &second, &third] { 3280 assert_accepted(relay.handle_event(event.clone()).expect("event"), event); 3281 } 3282 3283 let unbounded = filter_from_value(&serde_json::json!({ 3284 "authors": [first.unsigned().pubkey().as_str()], 3285 "kinds": [1] 3286 })) 3287 .expect("unbounded"); 3288 let client_limited = filter_from_value(&serde_json::json!({ 3289 "authors": [first.unsigned().pubkey().as_str()], 3290 "kinds": [1], 3291 "limit": 1 3292 })) 3293 .expect("client limited"); 3294 3295 assert_eq!( 3296 relay 3297 .handle_count_protocol( 3298 SubscriptionId::new("count-unbounded").expect("sub"), 3299 vec![unbounded] 3300 ) 3301 .expect("count"), 3302 RelayMessage::Count { 3303 subscription_id: SubscriptionId::new("count-unbounded").expect("sub"), 3304 count: 3, 3305 hll: None 3306 } 3307 ); 3308 assert_eq!( 3309 relay 3310 .handle_count_protocol( 3311 SubscriptionId::new("count-client-limited").expect("sub"), 3312 vec![client_limited] 3313 ) 3314 .expect("count"), 3315 RelayMessage::Count { 3316 subscription_id: SubscriptionId::new("count-client-limited").expect("sub"), 3317 count: 3, 3318 hll: None 3319 } 3320 ); 3321 } 3322 3323 #[test] 3324 fn base_relay_event_path_rejects_invalid_signatures_and_skips_ephemeral_storage() { 3325 let relay = test_relay("base-relay-event-store-path", 8); 3326 let valid = signed_public_event(7, 1, Vec::new(), "valid"); 3327 let signature_source = signed_public_event(8, 1, Vec::new(), "signature source"); 3328 let invalid = Event::new( 3329 valid.id().clone(), 3330 valid.unsigned().clone(), 3331 signature_source.sig().clone(), 3332 ); 3333 let ephemeral = signed_public_event(7, 20_001, Vec::new(), "ephemeral"); 3334 3335 assert!(matches!( 3336 relay.handle_event(invalid.clone()).expect("invalid"), 3337 RelayMessage::Ok { 3338 event_id, 3339 accepted: false, 3340 message 3341 } if event_id == *invalid.id() 3342 && message.starts_with("invalid:") 3343 )); 3344 assert_eq!(count_kind(&relay, 1), 0); 3345 3346 assert_accepted(relay.handle_event(valid.clone()).expect("valid"), &valid); 3347 assert_eq!( 3348 relay.handle_event(valid.clone()).expect("duplicate"), 3349 RelayMessage::Ok { 3350 event_id: valid.id().clone(), 3351 accepted: true, 3352 message: "duplicate: already have this event".to_owned() 3353 } 3354 ); 3355 assert_eq!(count_kind(&relay, 1), 1); 3356 3357 assert_accepted( 3358 relay.handle_event(ephemeral.clone()).expect("ephemeral"), 3359 &ephemeral, 3360 ); 3361 assert_eq!(count_kind(&relay, 20_001), 0); 3362 } 3363 3364 #[test] 3365 fn base_relay_pocket_event_path_preserves_event_admission_behavior() { 3366 let relay = test_relay("base-relay-pocket-event-store-path", 8); 3367 let tags = PocketOwnedTags::empty(); 3368 let protected_tags = PocketOwnedTags::new(&[["-"]]).expect("protected tags"); 3369 let valid_pocket = signed_pocket_event(7, 1, &tags, b"valid"); 3370 let signature_source = signed_pocket_event(8, 1, &tags, b"valid"); 3371 let invalid_pocket = PocketOwnedEvent::new( 3372 valid_pocket.id(), 3373 valid_pocket.kind(), 3374 valid_pocket.pubkey(), 3375 signature_source.sig(), 3376 valid_pocket.tags().expect("tags"), 3377 valid_pocket.created_at(), 3378 valid_pocket.content(), 3379 ) 3380 .expect("invalid pocket"); 3381 let ephemeral_pocket = signed_pocket_event(7, 20_001, &tags, b"ephemeral"); 3382 let protected_pocket = signed_pocket_event(7, 1, &protected_tags, b"protected"); 3383 3384 assert!( 3385 rejected_message(relay.handle_pocket_event(&invalid_pocket).expect("invalid")) 3386 .starts_with("invalid:") 3387 ); 3388 assert_eq!(count_kind(&relay, 1), 0); 3389 3390 assert_pocket_accepted( 3391 relay 3392 .handle_pocket_event(&valid_pocket) 3393 .expect("valid pocket"), 3394 &valid_pocket, 3395 ); 3396 assert_eq!( 3397 relay.handle_pocket_event(&valid_pocket).expect("duplicate"), 3398 RelayMessage::Ok { 3399 event_id: pocket_event_id(&valid_pocket), 3400 accepted: true, 3401 message: "duplicate: already have this event".to_owned() 3402 } 3403 ); 3404 assert_eq!(count_kind(&relay, 1), 1); 3405 3406 assert_pocket_accepted( 3407 relay 3408 .handle_pocket_event(&ephemeral_pocket) 3409 .expect("ephemeral"), 3410 &ephemeral_pocket, 3411 ); 3412 assert_eq!(count_kind(&relay, 20_001), 0); 3413 3414 assert_eq!( 3415 rejected_message( 3416 relay 3417 .handle_pocket_event(&protected_pocket) 3418 .expect("protected") 3419 ), 3420 "auth-required: protected event requires authenticated event author" 3421 ); 3422 assert_pocket_accepted( 3423 relay 3424 .handle_pocket_event_with_auth(&protected_pocket, &authenticated_state(7)) 3425 .expect("protected auth"), 3426 &protected_pocket, 3427 ); 3428 } 3429 3430 #[test] 3431 fn group_write_source_uses_atomic_service_boundary() { 3432 let core_source = include_str!("core.rs"); 3433 let group_source = include_str!("../groups.rs"); 3434 3435 assert!(core_source.contains("groups.store_group_pocket_event")); 3436 assert!(!core_source.contains(concat!("groups.", "check_event"))); 3437 assert!(!core_source.contains(concat!("groups.", "after_source_event_stored"))); 3438 assert!(!core_source.contains(concat!( 3439 "let tangle_event = ", 3440 "pocket_event_to_tangle(event)?;" 3441 ))); 3442 assert!(!group_source.contains("pub(crate) fn check_event(")); 3443 assert!(!group_source.contains("pub(crate) fn after_source_event_stored(")); 3444 } 3445 3446 #[test] 3447 fn base_relay_event_path_preserves_chorus_parity() { 3448 let owner = signer(7).public_key().clone(); 3449 let relay = test_relay_with_groups( 3450 "base-relay-event-chorus-parity", 3451 8, 3452 &enabled_groups_for_owner(&owner), 3453 ); 3454 let valid = signed_public_event(7, 1, Vec::new(), "valid"); 3455 let signature_source = signed_public_event(8, 1, Vec::new(), "signature source"); 3456 let invalid = Event::new( 3457 valid.id().clone(), 3458 valid.unsigned().clone(), 3459 signature_source.sig().clone(), 3460 ); 3461 let ephemeral = signed_public_event(7, 20_001, Vec::new(), "ephemeral"); 3462 let protected = signed_public_event( 3463 7, 3464 1, 3465 vec![Tag::from_parts("-", &[]).expect("protected")], 3466 "protected", 3467 ); 3468 let group_create = signed_group_create_event(7, "ParityFarm"); 3469 let empty_auth = BaseAuthState::new("wss://relay.radroots.test", 60, 600).expect("auth"); 3470 3471 assert!( 3472 rejected_message(relay.handle_event(invalid.clone()).expect("invalid")) 3473 .starts_with("invalid:") 3474 ); 3475 assert_eq!(count_kind(&relay, 1), 0); 3476 3477 assert_accepted(relay.handle_event(valid.clone()).expect("valid"), &valid); 3478 assert_eq!(count_kind(&relay, 1), 1); 3479 assert_eq!( 3480 relay.handle_event(valid.clone()).expect("duplicate"), 3481 RelayMessage::Ok { 3482 event_id: valid.id().clone(), 3483 accepted: true, 3484 message: "duplicate: already have this event".to_owned() 3485 } 3486 ); 3487 assert_eq!(count_kind(&relay, 1), 1); 3488 3489 assert_accepted( 3490 relay.handle_event(ephemeral.clone()).expect("ephemeral"), 3491 &ephemeral, 3492 ); 3493 assert_eq!(count_kind(&relay, 20_001), 0); 3494 3495 assert_eq!( 3496 rejected_message(relay.handle_event(protected.clone()).expect("protected")), 3497 "auth-required: protected event requires authenticated event author" 3498 ); 3499 assert_eq!( 3500 rejected_message( 3501 relay 3502 .handle_event_with_auth(group_create.clone(), &empty_auth) 3503 .expect("group unauth") 3504 ), 3505 "auth-required: group event author must authenticate with AUTH" 3506 ); 3507 assert_eq!(count_kind(&relay, KIND_GROUP_CREATE_GROUP), 0); 3508 3509 assert_accepted( 3510 relay 3511 .handle_event_with_auth(group_create.clone(), &authenticated_state(7)) 3512 .expect("group auth"), 3513 &group_create, 3514 ); 3515 assert_eq!(count_kind(&relay, KIND_GROUP_CREATE_GROUP), 1); 3516 assert!( 3517 relay 3518 .group_projection() 3519 .expect("projection") 3520 .group(&GroupId::new("ParityFarm").expect("group")) 3521 .is_some() 3522 ); 3523 } 3524 3525 #[test] 3526 fn base_relay_enforces_nip70_protected_event_author_auth() { 3527 let relay = test_relay("base-relay-nip70-protected", 8); 3528 let protected = signed_public_event( 3529 7, 3530 1, 3531 vec![Tag::from_parts("-", &[]).expect("protected")], 3532 "protected", 3533 ); 3534 3535 assert_eq!( 3536 rejected_message(relay.handle_event(protected.clone()).expect("unauth")), 3537 "auth-required: protected event requires authenticated event author" 3538 ); 3539 assert_eq!(count_kind(&relay, 1), 0); 3540 assert_eq!( 3541 rejected_message( 3542 relay 3543 .handle_event_with_auth(protected.clone(), &authenticated_state(8)) 3544 .expect("wrong auth") 3545 ), 3546 "auth-required: protected event requires authenticated event author" 3547 ); 3548 assert_eq!(count_kind(&relay, 1), 0); 3549 assert_accepted( 3550 relay 3551 .handle_event_with_auth(protected.clone(), &authenticated_state(7)) 3552 .expect("author auth"), 3553 &protected, 3554 ); 3555 assert_eq!(count_kind(&relay, 1), 1); 3556 } 3557 3558 #[test] 3559 fn base_relay_rejects_group_marked_events_before_group_service() { 3560 let relay = test_relay("base-relay-group-reject", 4); 3561 let event = signed_public_event( 3562 7, 3563 1, 3564 vec![Tag::from_parts("h", &["public-group"]).expect("group")], 3565 "hello", 3566 ); 3567 3568 assert_eq!( 3569 relay.handle_event(event.clone()).expect("event"), 3570 RelayMessage::Ok { 3571 event_id: event.id().clone(), 3572 accepted: false, 3573 message: "blocked: NIP-29 group events are not accepted before group service" 3574 .to_owned() 3575 } 3576 ); 3577 } 3578 3579 #[test] 3580 fn base_relay_rejects_client_submitted_relay_generated_group_state() { 3581 let relay = test_relay("base-relay-generated-group-reject", 4); 3582 for kind in NIP29_RELAY_GENERATED_KIND_VALUES { 3583 let event = signed_public_event( 3584 7, 3585 kind.into(), 3586 vec![Tag::from_parts("d", &["public-group"]).expect("group")], 3587 "", 3588 ); 3589 3590 assert_eq!( 3591 relay.handle_event(event.clone()).expect("event"), 3592 RelayMessage::Ok { 3593 event_id: event.id().clone(), 3594 accepted: false, 3595 message: 3596 "blocked: relay-generated group state events cannot be submitted by clients" 3597 .to_owned() 3598 } 3599 ); 3600 } 3601 } 3602 3603 #[test] 3604 fn base_relay_initializes_group_service_from_config() { 3605 let owner = signer(7).public_key().clone(); 3606 let relay = test_relay_with_groups( 3607 "base-relay-groups-enabled", 3608 4, 3609 &enabled_groups_for_owner(&owner), 3610 ); 3611 let disabled = test_relay_with_groups("base-relay-groups-disabled", 4, &disabled_groups()); 3612 3613 assert!(relay.groups_enabled()); 3614 assert_eq!( 3615 relay 3616 .readiness_state() 3617 .response() 3618 .checks 3619 .group_outbox_replay, 3620 "ready" 3621 ); 3622 assert!( 3623 relay 3624 .group_projection() 3625 .expect("projection") 3626 .groups() 3627 .is_empty() 3628 ); 3629 assert!(!disabled.groups_enabled()); 3630 assert_eq!( 3631 disabled 3632 .readiness_state() 3633 .response() 3634 .checks 3635 .group_outbox_replay, 3636 "ready" 3637 ); 3638 assert!(disabled.group_projection().is_none()); 3639 } 3640 3641 #[test] 3642 fn group_event_write_requires_auth_before_storage() { 3643 let owner = signer(7).public_key().clone(); 3644 let relay = test_relay_with_groups( 3645 "base-relay-group-auth-required", 3646 4, 3647 &enabled_groups_for_owner(&owner), 3648 ); 3649 let auth = BaseAuthState::new("wss://relay.radroots.test", 60, 600).expect("auth"); 3650 let event = signed_group_create_event(7, "Farm"); 3651 3652 assert_eq!( 3653 relay 3654 .handle_event_with_auth(event.clone(), &auth) 3655 .expect("event"), 3656 RelayMessage::Ok { 3657 event_id: event.id().clone(), 3658 accepted: false, 3659 message: "auth-required: group event author must authenticate with AUTH".to_owned() 3660 } 3661 ); 3662 assert!( 3663 relay 3664 .group_projection() 3665 .expect("projection") 3666 .group(&GroupId::new("Farm").expect("group")) 3667 .is_none() 3668 ); 3669 assert_eq!(count_kind(&relay, KIND_GROUP_CREATE_GROUP), 0); 3670 } 3671 3672 #[test] 3673 fn group_create_updates_projection_and_stores_generated_snapshots() { 3674 let owner = signer(7).public_key().clone(); 3675 let relay = test_relay_with_groups( 3676 "base-relay-group-create", 3677 4, 3678 &enabled_groups_for_owner(&owner), 3679 ); 3680 let auth = authenticated_state(7); 3681 let event = signed_group_create_event(7, "Farm"); 3682 3683 assert_eq!( 3684 relay 3685 .handle_event_with_auth(event.clone(), &auth) 3686 .expect("event"), 3687 RelayMessage::Ok { 3688 event_id: event.id().clone(), 3689 accepted: true, 3690 message: String::new() 3691 } 3692 ); 3693 3694 let group_id = GroupId::new("Farm").expect("group"); 3695 assert!( 3696 relay 3697 .group_projection() 3698 .expect("projection") 3699 .group(&group_id) 3700 .is_some() 3701 ); 3702 assert_eq!(count_kind(&relay, KIND_GROUP_CREATE_GROUP), 1); 3703 assert_eq!(count_kind(&relay, KIND_GROUP_METADATA), 1); 3704 assert_eq!(count_kind(&relay, KIND_GROUP_ADMINS), 1); 3705 } 3706 3707 #[test] 3708 fn group_join_materializes_relay_membership_event() { 3709 let owner = signer(7).public_key().clone(); 3710 let joiner = signer(8).public_key().clone(); 3711 let relay = test_relay_with_groups( 3712 "base-relay-group-join", 3713 4, 3714 &enabled_groups_for_owner_with_public_join(&owner), 3715 ); 3716 let create = signed_group_create_event(7, "Farm"); 3717 assert_accepted( 3718 relay 3719 .handle_event_with_auth(create.clone(), &authenticated_state(7)) 3720 .expect("create"), 3721 &create, 3722 ); 3723 let join = signed_event_at( 3724 8, 3725 KIND_GROUP_JOIN_REQUEST.into(), 3726 vec![Tag::from_parts("h", &["Farm"]).expect("h")], 3727 "", 3728 1_714_124_434, 3729 ); 3730 3731 assert_eq!( 3732 relay 3733 .handle_event_with_auth(join.clone(), &authenticated_state(8)) 3734 .expect("join"), 3735 RelayMessage::Ok { 3736 event_id: join.id().clone(), 3737 accepted: true, 3738 message: String::new() 3739 } 3740 ); 3741 3742 assert_eq!(count_kind(&relay, KIND_GROUP_PUT_USER), 1); 3743 assert_eq!( 3744 relay 3745 .group_projection() 3746 .expect("projection") 3747 .member(&GroupId::new("Farm").expect("group"), &joiner) 3748 .expect("member") 3749 .status(), 3750 MemberStatus::Member 3751 ); 3752 } 3753 3754 #[test] 3755 fn group_join_requires_public_join_policy() { 3756 let owner = signer(7).public_key().clone(); 3757 let relay = test_relay_with_groups( 3758 "base-relay-group-join-default-deny", 3759 4, 3760 &enabled_groups_for_owner(&owner), 3761 ); 3762 let create = signed_group_create_event(7, "Farm"); 3763 relay 3764 .handle_event_with_auth(create, &authenticated_state(7)) 3765 .expect("create"); 3766 let join = signed_event_at( 3767 8, 3768 KIND_GROUP_JOIN_REQUEST.into(), 3769 vec![Tag::from_parts("h", &["Farm"]).expect("h")], 3770 "", 3771 1_714_124_434, 3772 ); 3773 3774 assert_eq!( 3775 rejected_message( 3776 relay 3777 .handle_event_with_auth(join, &authenticated_state(8)) 3778 .expect("join") 3779 ), 3780 "restricted: group is unavailable" 3781 ); 3782 assert_eq!(count_kind(&relay, KIND_GROUP_PUT_USER), 0); 3783 } 3784 3785 #[test] 3786 fn group_metadata_edit_replaces_generated_metadata_snapshot() { 3787 let owner = signer(7).public_key().clone(); 3788 let mut relay = test_relay_with_groups( 3789 "base-relay-group-metadata-edit", 3790 4, 3791 &enabled_groups_for_owner(&owner), 3792 ); 3793 let auth = authenticated_state(7); 3794 let create = signed_group_create_event(7, "Farm"); 3795 assert_accepted( 3796 relay 3797 .handle_event_with_auth(create.clone(), &auth) 3798 .expect("create"), 3799 &create, 3800 ); 3801 let edit = signed_event_at( 3802 7, 3803 KIND_GROUP_EDIT_METADATA.into(), 3804 vec![h("Farm"), name("Market")], 3805 "", 3806 1_714_124_436, 3807 ); 3808 assert_accepted( 3809 relay 3810 .handle_event_with_auth(edit.clone(), &auth) 3811 .expect("edit"), 3812 &edit, 3813 ); 3814 3815 let group_id = GroupId::new("Farm").expect("group"); 3816 { 3817 let projection = relay.group_projection().expect("projection"); 3818 let group = projection.group(&group_id).expect("group"); 3819 assert_eq!(group.metadata().name(), Some("Market")); 3820 } 3821 let metadata = query_filter( 3822 &mut relay, 3823 "metadata-edit", 3824 filter_group_tag(KIND_GROUP_METADATA, "d", "Farm"), 3825 ); 3826 assert_eq!(metadata.len(), 1); 3827 assert!(has_tag(&metadata[0], "d", &["Farm"])); 3828 assert!(has_tag(&metadata[0], "name", &["Market"])); 3829 assert_eq!(count_kind(&relay, KIND_GROUP_METADATA), 1); 3830 } 3831 3832 #[test] 3833 fn group_member_moderation_join_leave_and_snapshots_flow() { 3834 let owner = signer(7).public_key().clone(); 3835 let member = signer(8).public_key().clone(); 3836 let target = signer(9).public_key().clone(); 3837 let relay = test_relay_with_groups( 3838 "base-relay-group-member-flow", 3839 4, 3840 &enabled_groups_for_owner_with_public_join(&owner), 3841 ); 3842 let owner_auth = authenticated_state(7); 3843 let member_auth = authenticated_state(8); 3844 let target_auth = authenticated_state(9); 3845 relay 3846 .handle_event_with_auth(signed_group_create_event(7, "Farm"), &owner_auth) 3847 .expect("create"); 3848 let rejected_add = signed_event_at( 3849 9, 3850 KIND_GROUP_PUT_USER.into(), 3851 vec![h("Farm"), p(&target)], 3852 "", 3853 1_714_124_434, 3854 ); 3855 assert_eq!( 3856 rejected_message( 3857 relay 3858 .handle_event_with_auth(rejected_add.clone(), &target_auth) 3859 .expect("rejected add") 3860 ), 3861 "restricted: missing group capability manage_members" 3862 ); 3863 let add = signed_event_at( 3864 7, 3865 KIND_GROUP_PUT_USER.into(), 3866 vec![h("Farm"), p(&member)], 3867 "", 3868 1_714_124_435, 3869 ); 3870 assert_accepted( 3871 relay 3872 .handle_event_with_auth(add.clone(), &owner_auth) 3873 .expect("add"), 3874 &add, 3875 ); 3876 assert_member_status(&relay, "Farm", &member, MemberStatus::Member); 3877 assert_eq!(count_kind(&relay, KIND_GROUP_MEMBERS), 1); 3878 3879 let remove = signed_event_at( 3880 7, 3881 KIND_GROUP_REMOVE_USER.into(), 3882 vec![h("Farm"), p(&member)], 3883 "", 3884 1_714_124_436, 3885 ); 3886 assert_accepted( 3887 relay 3888 .handle_event_with_auth(remove.clone(), &owner_auth) 3889 .expect("remove"), 3890 &remove, 3891 ); 3892 assert_member_status(&relay, "Farm", &member, MemberStatus::Removed); 3893 assert_eq!(count_kind(&relay, KIND_GROUP_MEMBERS), 1); 3894 3895 let join = signed_event_at( 3896 8, 3897 KIND_GROUP_JOIN_REQUEST.into(), 3898 vec![h("Farm")], 3899 "", 3900 1_714_124_437, 3901 ); 3902 assert_accepted( 3903 relay 3904 .handle_event_with_auth(join.clone(), &member_auth) 3905 .expect("join"), 3906 &join, 3907 ); 3908 assert_member_status(&relay, "Farm", &member, MemberStatus::Member); 3909 let duplicate_join = signed_event_at( 3910 8, 3911 KIND_GROUP_JOIN_REQUEST.into(), 3912 vec![h("Farm")], 3913 "", 3914 1_714_124_438, 3915 ); 3916 assert_eq!( 3917 rejected_message( 3918 relay 3919 .handle_event_with_auth(duplicate_join, &member_auth) 3920 .expect("duplicate join") 3921 ), 3922 "duplicate: group member already exists" 3923 ); 3924 3925 let leave = signed_event_at( 3926 8, 3927 KIND_GROUP_LEAVE_REQUEST.into(), 3928 vec![h("Farm")], 3929 "", 3930 1_714_124_439, 3931 ); 3932 assert_accepted( 3933 relay 3934 .handle_event_with_auth(leave.clone(), &member_auth) 3935 .expect("leave"), 3936 &leave, 3937 ); 3938 assert_member_status(&relay, "Farm", &member, MemberStatus::Removed); 3939 assert_eq!(count_kind(&relay, KIND_GROUP_REMOVE_USER), 2); 3940 let duplicate_leave = signed_event_at( 3941 8, 3942 KIND_GROUP_LEAVE_REQUEST.into(), 3943 vec![h("Farm")], 3944 "", 3945 1_714_124_440, 3946 ); 3947 assert_eq!( 3948 rejected_message( 3949 relay 3950 .handle_event_with_auth(duplicate_leave, &member_auth) 3951 .expect("duplicate leave") 3952 ), 3953 "duplicate: group member does not exist" 3954 ); 3955 } 3956 3957 #[test] 3958 fn group_delete_event_moderation_hides_target_and_validates_group() { 3959 let owner = signer(7).public_key().clone(); 3960 let outsider_auth = authenticated_state(8); 3961 let owner_auth = authenticated_state(7); 3962 let relay = test_relay_with_groups( 3963 "base-relay-group-delete-event", 3964 4, 3965 &enabled_groups_for_owner(&owner), 3966 ); 3967 relay 3968 .handle_event_with_auth(signed_group_create_event(7, "Farm"), &owner_auth) 3969 .expect("create farm"); 3970 relay 3971 .handle_event_with_auth(signed_group_create_event(7, "Other"), &owner_auth) 3972 .expect("create other"); 3973 let target = signed_event_at(7, 1, vec![h("Farm")], "harvest", 1_714_124_434); 3974 let other = signed_event_at(7, 1, vec![h("Other")], "other", 1_714_124_435); 3975 relay 3976 .handle_event_with_auth(target.clone(), &owner_auth) 3977 .expect("target"); 3978 relay 3979 .handle_event_with_auth(other.clone(), &owner_auth) 3980 .expect("other"); 3981 3982 let wrong_group = signed_event_at( 3983 7, 3984 KIND_GROUP_DELETE_EVENT.into(), 3985 vec![h("Farm"), e(other.id())], 3986 "", 3987 1_714_124_436, 3988 ); 3989 assert_eq!( 3990 rejected_message( 3991 relay 3992 .handle_event_with_auth(wrong_group, &owner_auth) 3993 .expect("wrong group") 3994 ), 3995 "invalid: delete target event is not in group" 3996 ); 3997 let unauthorized = signed_event_at( 3998 8, 3999 KIND_GROUP_DELETE_EVENT.into(), 4000 vec![h("Farm"), e(target.id())], 4001 "", 4002 1_714_124_437, 4003 ); 4004 assert_eq!( 4005 rejected_message( 4006 relay 4007 .handle_event_with_auth(unauthorized, &outsider_auth) 4008 .expect("unauthorized") 4009 ), 4010 "restricted: missing group capability delete_events" 4011 ); 4012 assert_eq!( 4013 count_filter( 4014 &relay, 4015 "target-before-delete", 4016 filter_group_tag(1, "h", "Farm") 4017 ), 4018 1 4019 ); 4020 4021 let delete = signed_event_at( 4022 7, 4023 KIND_GROUP_DELETE_EVENT.into(), 4024 vec![h("Farm"), e(target.id())], 4025 "", 4026 1_714_124_438, 4027 ); 4028 assert_accepted( 4029 relay 4030 .handle_event_with_auth(delete.clone(), &owner_auth) 4031 .expect("delete"), 4032 &delete, 4033 ); 4034 4035 assert_eq!( 4036 count_filter( 4037 &relay, 4038 "target-after-delete", 4039 filter_group_tag(1, "h", "Farm") 4040 ), 4041 0 4042 ); 4043 assert_eq!( 4044 count_filter( 4045 &relay, 4046 "delete-event-marker", 4047 filter_group_tag(KIND_GROUP_DELETE_EVENT, "h", "Farm") 4048 ), 4049 1 4050 ); 4051 } 4052 4053 #[test] 4054 fn group_delete_tombstone_hides_events_and_rejects_future_writes() { 4055 let owner = signer(7).public_key().clone(); 4056 let auth = authenticated_state(7); 4057 let relay = test_relay_with_groups( 4058 "base-relay-group-delete-tombstone", 4059 4, 4060 &enabled_groups_for_owner(&owner), 4061 ); 4062 relay 4063 .handle_event_with_auth(signed_group_create_event(7, "Farm"), &auth) 4064 .expect("create"); 4065 let normal = signed_event_at(7, 1, vec![h("Farm")], "harvest", 1_714_124_434); 4066 relay.handle_event_with_auth(normal, &auth).expect("normal"); 4067 let delete_group = signed_event_at( 4068 7, 4069 KIND_GROUP_DELETE_GROUP.into(), 4070 vec![h("Farm")], 4071 "", 4072 1_714_124_435, 4073 ); 4074 assert_accepted( 4075 relay 4076 .handle_event_with_auth(delete_group.clone(), &auth) 4077 .expect("delete group"), 4078 &delete_group, 4079 ); 4080 4081 let future = signed_event_at(7, 1, vec![h("Farm")], "future", 1_714_124_436); 4082 assert_eq!( 4083 rejected_message(relay.handle_event_with_auth(future, &auth).expect("future")), 4084 "blocked: group is deleted" 4085 ); 4086 assert_eq!( 4087 count_filter( 4088 &relay, 4089 "deleted-group-normal", 4090 filter_group_tag(1, "h", "Farm") 4091 ), 4092 0 4093 ); 4094 assert_eq!( 4095 count_filter( 4096 &relay, 4097 "deleted-group-marker", 4098 filter_group_tag(KIND_GROUP_DELETE_GROUP, "h", "Farm") 4099 ), 4100 1 4101 ); 4102 } 4103 4104 #[test] 4105 fn strict_closed_restricted_hidden_and_disabled_invite_flows() { 4106 let owner = signer(7).public_key().clone(); 4107 let outsider_auth = authenticated_state(8); 4108 let owner_auth = authenticated_state(7); 4109 let relay = test_relay_with_groups( 4110 "base-relay-group-strict-policy-flow", 4111 4, 4112 &enabled_groups_for_owner(&owner), 4113 ); 4114 relay 4115 .handle_event_with_auth( 4116 signed_group_create_event_with_tags(7, "Restricted", vec![restricted()], 1), 4117 &owner_auth, 4118 ) 4119 .expect("restricted create"); 4120 let restricted_write = 4121 signed_event_at(8, 1, vec![h("Restricted")], "restricted", 1_714_124_434); 4122 assert_eq!( 4123 rejected_message( 4124 relay 4125 .handle_event_with_auth(restricted_write, &outsider_auth) 4126 .expect("restricted write") 4127 ), 4128 "restricted: group is unavailable" 4129 ); 4130 4131 relay 4132 .handle_event_with_auth( 4133 signed_group_create_event_with_tags(7, "Closed", vec![closed()], 2), 4134 &owner_auth, 4135 ) 4136 .expect("closed create"); 4137 let closed_join = signed_event_at( 4138 8, 4139 KIND_GROUP_JOIN_REQUEST.into(), 4140 vec![h("Closed")], 4141 "", 4142 1_714_124_435, 4143 ); 4144 assert_eq!( 4145 rejected_message( 4146 relay 4147 .handle_event_with_auth(closed_join, &outsider_auth) 4148 .expect("closed join") 4149 ), 4150 "restricted: group is unavailable" 4151 ); 4152 let closed_normal = signed_event_at(8, 1, vec![h("Closed")], "open", 1_714_124_436); 4153 assert_accepted( 4154 relay 4155 .handle_event_with_auth(closed_normal.clone(), &outsider_auth) 4156 .expect("closed normal"), 4157 &closed_normal, 4158 ); 4159 4160 relay 4161 .handle_event_with_auth( 4162 signed_group_create_event_with_tags(7, "Hidden", vec![hidden()], 3), 4163 &owner_auth, 4164 ) 4165 .expect("hidden create"); 4166 assert_eq!( 4167 count_filter( 4168 &relay, 4169 "hidden-unauth", 4170 filter_group_tag(KIND_GROUP_METADATA, "d", "Hidden") 4171 ), 4172 0 4173 ); 4174 assert_eq!( 4175 count_filter_with_auth( 4176 &relay, 4177 "hidden-owner", 4178 filter_group_tag(KIND_GROUP_METADATA, "d", "Hidden"), 4179 &owner_auth 4180 ), 4181 1 4182 ); 4183 4184 let invite = signed_event_at( 4185 7, 4186 KIND_GROUP_CREATE_INVITE.into(), 4187 vec![h("Closed")], 4188 "", 4189 1_714_124_437, 4190 ); 4191 assert_eq!( 4192 rejected_message( 4193 relay 4194 .handle_event_with_auth(invite, &owner_auth) 4195 .expect("invite") 4196 ), 4197 "restricted: invites not enabled" 4198 ); 4199 } 4200 4201 #[test] 4202 fn private_group_req_and_count_use_reader_auth() { 4203 let owner = signer(7).public_key().clone(); 4204 let auth = authenticated_state(7); 4205 let mut relay = test_relay_with_groups( 4206 "base-relay-private-read", 4207 4, 4208 &enabled_groups_for_owner(&owner), 4209 ); 4210 relay 4211 .handle_event_with_auth(signed_private_group_create_event(7, "Farm"), &auth) 4212 .expect("create"); 4213 let private_event = signed_event_at( 4214 7, 4215 1, 4216 vec![Tag::from_parts("h", &["Farm"]).expect("h")], 4217 "private harvest", 4218 1_714_124_435, 4219 ); 4220 relay 4221 .handle_event_with_auth(private_event.clone(), &auth) 4222 .expect("private event"); 4223 4224 let unauth_sub = SubscriptionId::new("private-unauth").expect("sub"); 4225 let auth_sub = SubscriptionId::new("private-auth").expect("sub"); 4226 assert_eq!( 4227 relay 4228 .handle_protocol_req_for_test(unauth_sub.clone(), vec![filter_kind(1)]) 4229 .expect("unauth req"), 4230 vec![RelayMessage::Closed { 4231 subscription_id: unauth_sub, 4232 message: "auth-required: authentication required to read group events".to_owned() 4233 }] 4234 ); 4235 assert_eq!(relay.active_subscription_count(), 0); 4236 assert!(matches!( 4237 relay 4238 .handle_protocol_req_with_auth_for_test(auth_sub.clone(), vec![filter_kind(1)], &auth) 4239 .expect("auth req") 4240 .as_slice(), 4241 [RelayMessage::Event { subscription_id, event }, RelayMessage::Eose(eose)] 4242 if subscription_id == &auth_sub && event.id() == private_event.id() && eose == &auth_sub 4243 )); 4244 assert_eq!(count_kind(&relay, 1), 0); 4245 assert_eq!(count_kind_with_auth(&relay, 1, &auth), 1); 4246 assert_eq!(count_kind(&relay, KIND_GROUP_METADATA), 1); 4247 assert_eq!(count_kind(&relay, KIND_GROUP_ADMINS), 1); 4248 assert_eq!(count_kind(&relay, KIND_GROUP_MEMBERS), 0); 4249 assert_eq!(count_kind_with_auth(&relay, KIND_GROUP_METADATA, &auth), 1); 4250 assert_eq!(count_kind_with_auth(&relay, KIND_GROUP_ADMINS, &auth), 1); 4251 } 4252 4253 #[test] 4254 fn private_and_hidden_group_offset_lookup_uses_reader_auth() { 4255 let owner = signer(7).public_key().clone(); 4256 let owner_auth = authenticated_state(7); 4257 let unauth = BaseAuthState::new("wss://relay.radroots.test", 60, 600).expect("auth state"); 4258 let relay = test_relay_with_groups( 4259 "base-relay-private-offset-read", 4260 4, 4261 &enabled_groups_for_owner(&owner), 4262 ); 4263 relay 4264 .handle_event_with_auth(signed_private_group_create_event(7, "Farm"), &owner_auth) 4265 .expect("create"); 4266 let private_event = signed_event_at( 4267 7, 4268 1, 4269 vec![Tag::from_parts("h", &["Farm"]).expect("h")], 4270 "private harvest", 4271 1_714_124_435, 4272 ); 4273 let pocket = tangle_event_to_pocket(&private_event).expect("pocket"); 4274 let offset = StoreOffset::new(relay.store.store_event(&pocket).expect("store")); 4275 4276 assert_eq!( 4277 relay 4278 .event_by_offset_with_auth(offset, &unauth) 4279 .expect("unauth offset"), 4280 None 4281 ); 4282 let visible = relay 4283 .event_by_offset_with_auth(offset, &owner_auth) 4284 .expect("owner offset") 4285 .expect("visible"); 4286 let visible: &PocketEvent = &visible; 4287 assert_eq!(visible.id().as_hex_string(), private_event.id().as_str()); 4288 4289 relay 4290 .handle_event_with_auth( 4291 signed_group_create_event_with_tags(7, "HiddenFarm", vec![hidden()], 1_714_124_436), 4292 &owner_auth, 4293 ) 4294 .expect("hidden create"); 4295 let hidden_event = signed_event_at( 4296 7, 4297 1, 4298 vec![Tag::from_parts("h", &["HiddenFarm"]).expect("h")], 4299 "hidden harvest", 4300 1_714_124_437, 4301 ); 4302 let pocket = tangle_event_to_pocket(&hidden_event).expect("hidden pocket"); 4303 let offset = StoreOffset::new(relay.store.store_event(&pocket).expect("store hidden")); 4304 4305 assert_eq!( 4306 relay 4307 .event_by_offset_with_auth(offset, &unauth) 4308 .expect("hidden unauth offset"), 4309 None 4310 ); 4311 let visible = relay 4312 .event_by_offset_with_auth(offset, &owner_auth) 4313 .expect("hidden owner offset") 4314 .expect("hidden visible"); 4315 let visible: &PocketEvent = &visible; 4316 assert_eq!(visible.id().as_hex_string(), hidden_event.id().as_str()); 4317 } 4318 4319 #[test] 4320 fn private_group_live_fanout_uses_current_auth() { 4321 let owner = signer(7).public_key().clone(); 4322 let auth = authenticated_state(7); 4323 let mut relay = test_relay_with_groups( 4324 "base-relay-private-fanout", 4325 4, 4326 &enabled_groups_for_owner(&owner), 4327 ); 4328 relay 4329 .handle_event_with_auth(signed_private_group_create_event(7, "Farm"), &auth) 4330 .expect("create"); 4331 let subscription_id = SubscriptionId::new("fanout-current-auth").expect("sub"); 4332 relay 4333 .handle_protocol_req_for_test(subscription_id.clone(), vec![filter_kind(1)]) 4334 .expect("sub"); 4335 let private_event = signed_event_at( 4336 7, 4337 1, 4338 vec![Tag::from_parts("h", &["Farm"]).expect("h")], 4339 "private harvest", 4340 1_714_124_435, 4341 ); 4342 relay 4343 .handle_event_with_auth(private_event.clone(), &auth) 4344 .expect("private event"); 4345 4346 assert!(relay.fanout_protocol_for_test(&private_event).is_empty()); 4347 assert!(matches!( 4348 relay 4349 .fanout_protocol_with_group_auth_for_test( 4350 &private_event, 4351 &GroupAuthContext::new([owner]) 4352 ) 4353 .as_slice(), 4354 [RelayMessage::Event { 4355 subscription_id: delivered, 4356 event 4357 }] if delivered == &subscription_id && event.id() == private_event.id() 4358 )); 4359 } 4360 4361 #[test] 4362 fn live_subscription_delivery_volume_does_not_close_subscription() { 4363 let mut relay = test_relay("base-relay-delivery-volume", 1); 4364 let subscription_id = SubscriptionId::new("sub-volume").expect("sub"); 4365 let filter = filter_from_value(&serde_json::json!({"kinds":[1]})).expect("filter"); 4366 relay 4367 .handle_protocol_req_for_test(subscription_id.clone(), vec![filter]) 4368 .expect("req"); 4369 let first = signed_public_event(7, 1, Vec::new(), "first"); 4370 let second = signed_public_event(7, 1, Vec::new(), "second"); 4371 4372 assert!(matches!( 4373 relay.fanout_protocol_for_test(&first).as_slice(), 4374 [RelayMessage::Event { .. }] 4375 )); 4376 assert!(matches!( 4377 relay.fanout_protocol_for_test(&second).as_slice(), 4378 [RelayMessage::Event { .. }] 4379 )); 4380 assert_eq!(relay.active_subscription_count(), 1); 4381 } 4382 4383 #[test] 4384 fn base_relay_shutdown_closes_live_subscriptions_and_syncs_store() { 4385 let config = test_store_config("base-relay-shutdown"); 4386 let mut relay = 4387 BaseRelay::open(&config, relay_limits(4), PocketQueryConfig::default()).expect("relay"); 4388 let event = signed_public_event(7, 1, Vec::new(), "shutdown"); 4389 let subscription_id = SubscriptionId::new("sub-shutdown").expect("sub"); 4390 4391 assert_accepted(relay.handle_event(event.clone()).expect("event"), &event); 4392 relay 4393 .handle_protocol_req_for_test(subscription_id, vec![filter_kind(1)]) 4394 .expect("req"); 4395 4396 assert_eq!(relay.active_subscription_count(), 1); 4397 4398 let report = relay.shutdown().expect("shutdown"); 4399 4400 assert_eq!(report.closed_subscriptions(), 1); 4401 assert_eq!(relay.active_subscription_count(), 0); 4402 assert!(relay.fanout_protocol_for_test(&event).is_empty()); 4403 4404 let reopened = BaseRelay::open(&config, relay_limits(4), PocketQueryConfig::default()) 4405 .expect("reopened"); 4406 assert_eq!(count_kind(&reopened, 1), 1); 4407 } 4408 4409 #[test] 4410 fn base_relay_client_message_dispatch_handles_count_and_auth() { 4411 let mut relay = test_relay("base-relay-dispatch", 4); 4412 let mut auth = 4413 BaseAuthState::new("wss://relay.radroots.test", 60, 600).expect("auth state"); 4414 auth.issue_challenge("challenge-a", UnixTimestamp::new(100)) 4415 .expect("challenge"); 4416 let auth_event = signed_auth_event(7, "challenge-a", 120); 4417 let count_id = SubscriptionId::new("count-a").expect("sub"); 4418 4419 assert_eq!( 4420 relay 4421 .handle_client_message( 4422 ClientMessage::Auth(auth_event.clone()), 4423 &mut auth, 4424 UnixTimestamp::new(120) 4425 ) 4426 .expect("auth"), 4427 vec![RelayMessage::Ok { 4428 event_id: auth_event.id().clone(), 4429 accepted: true, 4430 message: String::new() 4431 }] 4432 ); 4433 assert_eq!( 4434 relay 4435 .handle_client_message( 4436 ClientMessage::Count { 4437 subscription_id: count_id.clone(), 4438 filters: vec![Filter::empty()] 4439 }, 4440 &mut auth, 4441 UnixTimestamp::new(130) 4442 ) 4443 .expect("count"), 4444 vec![RelayMessage::Count { 4445 subscription_id: count_id, 4446 count: 0, 4447 hll: None 4448 }] 4449 ); 4450 } 4451 4452 #[test] 4453 fn base_relay_enforces_event_and_filter_runtime_limits() { 4454 let config = test_store_config("base-relay-event-filter-runtime-limits"); 4455 let mut relay = 4456 BaseRelay::open(&config, strict_relay_limits(), PocketQueryConfig::default()) 4457 .expect("relay"); 4458 let first = signed_public_event(7, 1, Vec::new(), "a"); 4459 let second = signed_event_at(8, 1, Vec::new(), "b", 1_714_124_434); 4460 4461 assert_accepted(relay.handle_event(first.clone()).expect("first"), &first); 4462 assert_accepted(relay.handle_event(second.clone()).expect("second"), &second); 4463 assert_eq!( 4464 rejected_message( 4465 relay 4466 .handle_event(signed_public_event(7, 1, Vec::new(), "abcde")) 4467 .expect("content") 4468 ), 4469 "invalid: event content length exceeds runtime max_content_length 4" 4470 ); 4471 assert_eq!( 4472 rejected_message( 4473 relay 4474 .handle_event(signed_public_event( 4475 7, 4476 1, 4477 vec![ 4478 Tag::from_parts("t", &["one"]).expect("tag"), 4479 Tag::from_parts("r", &["two"]).expect("tag"), 4480 ], 4481 "", 4482 )) 4483 .expect("tags") 4484 ), 4485 "invalid: event tag count exceeds runtime max_event_tags 1" 4486 ); 4487 assert_eq!( 4488 relay 4489 .handle_protocol_req_for_test( 4490 SubscriptionId::new("a").expect("sub"), 4491 vec![Filter::empty()] 4492 ) 4493 .expect("default limit") 4494 .len(), 4495 2 4496 ); 4497 assert!( 4498 relay 4499 .handle_protocol_req_for_test( 4500 SubscriptionId::new("a").expect("sub"), 4501 vec![Filter::empty(), Filter::empty()], 4502 ) 4503 .expect_err("filter count") 4504 .prefixed_message() 4505 .contains("max_filters_per_request 1") 4506 ); 4507 assert!( 4508 relay 4509 .handle_count_protocol( 4510 SubscriptionId::new("a").expect("sub"), 4511 vec![ 4512 filter_from_value(&serde_json::json!({"#t":["one", "two"]})) 4513 .expect("filter"), 4514 ], 4515 ) 4516 .expect_err("tag values") 4517 .prefixed_message() 4518 .contains("max_tag_values_per_filter 1") 4519 ); 4520 assert!( 4521 relay 4522 .handle_protocol_req_for_test( 4523 SubscriptionId::new("a").expect("sub"), 4524 vec![filter_from_value(&serde_json::json!({"limit": 3})).expect("filter")], 4525 ) 4526 .expect_err("max limit") 4527 .prefixed_message() 4528 .contains("max_limit 2") 4529 ); 4530 } 4531 4532 #[test] 4533 fn base_relay_enforces_subscription_id_and_count_limits() { 4534 let config = test_store_config("base-relay-subscription-limits"); 4535 let mut relay = 4536 BaseRelay::open(&config, strict_relay_limits(), PocketQueryConfig::default()) 4537 .expect("relay"); 4538 4539 assert!( 4540 relay 4541 .handle_protocol_req_for_test( 4542 SubscriptionId::new("abcde").expect("sub"), 4543 vec![Filter::empty()], 4544 ) 4545 .expect_err("sub id length") 4546 .prefixed_message() 4547 .contains("max_subid_length 4") 4548 ); 4549 relay 4550 .handle_protocol_req_for_test( 4551 SubscriptionId::new("a").expect("sub"), 4552 vec![Filter::empty()], 4553 ) 4554 .expect("first subscription"); 4555 assert!( 4556 relay 4557 .handle_protocol_req_for_test( 4558 SubscriptionId::new("b").expect("sub"), 4559 vec![Filter::empty()] 4560 ) 4561 .expect_err("subscription count") 4562 .prefixed_message() 4563 .contains("connection subscription limit exceeded") 4564 ); 4565 relay 4566 .handle_protocol_req_for_test( 4567 SubscriptionId::new("a").expect("sub"), 4568 vec![Filter::empty()], 4569 ) 4570 .expect("replace subscription"); 4571 } 4572 4573 fn test_relay(name: &str, max_pending_events: usize) -> BaseRelay { 4574 let config = test_store_config(name); 4575 BaseRelay::open( 4576 &config, 4577 relay_limits(max_pending_events), 4578 PocketQueryConfig::default(), 4579 ) 4580 .expect("relay") 4581 } 4582 4583 fn test_relay_with_groups( 4584 name: &str, 4585 max_pending_events: usize, 4586 groups: &tangle_groups::GroupRuntimeConfig, 4587 ) -> BaseRelay { 4588 let config = test_store_config(name); 4589 BaseRelay::open_with_groups( 4590 &config, 4591 relay_limits(max_pending_events), 4592 groups, 4593 PocketQueryConfig::default(), 4594 ) 4595 .expect("relay") 4596 } 4597 4598 fn relay_limits(max_pending_events: usize) -> BaseRelayLimits { 4599 BaseRelayLimits::new(BaseRelayLimitSettings { 4600 max_pending_events, 4601 max_subscription_id_length: 64, 4602 max_subscriptions: 64, 4603 max_filters_per_request: 10, 4604 max_tag_values_per_filter: 100, 4605 max_query_complexity: 610, 4606 max_event_tags: 200, 4607 max_content_length: 65_536, 4608 max_limit: 500, 4609 default_limit: 100, 4610 }) 4611 .expect("limits") 4612 } 4613 4614 fn strict_relay_limits() -> BaseRelayLimits { 4615 BaseRelayLimits::new(BaseRelayLimitSettings { 4616 max_pending_events: 4, 4617 max_subscription_id_length: 4, 4618 max_subscriptions: 1, 4619 max_filters_per_request: 1, 4620 max_tag_values_per_filter: 1, 4621 max_query_complexity: 4, 4622 max_event_tags: 1, 4623 max_content_length: 4, 4624 max_limit: 2, 4625 default_limit: 1, 4626 }) 4627 .expect("limits") 4628 } 4629 4630 fn test_store_config(name: &str) -> PocketStoreConfig { 4631 let root = std::env::temp_dir().join(format!("tangle-{name}-{}", std::process::id())); 4632 let _ = std::fs::remove_dir_all(&root); 4633 PocketStoreConfig::new(root.join("pocket"), PocketSyncPolicy::FlushOnShutdown) 4634 .expect("config") 4635 } 4636 4637 fn enabled_groups_for_owner(owner: &PublicKeyHex) -> tangle_groups::GroupRuntimeConfig { 4638 parse_group_runtime_config_json(&format!( 4639 r#"{{ 4640 "enabled": true, 4641 "canonical_relay_url": "wss://relay.radroots.test", 4642 "relay_secret": "{}", 4643 "owner_pubkeys": ["{}"] 4644 }}"#, 4645 "7".repeat(64), 4646 owner.as_str() 4647 )) 4648 .expect("groups") 4649 } 4650 4651 fn enabled_groups_for_owner_with_public_join( 4652 owner: &PublicKeyHex, 4653 ) -> tangle_groups::GroupRuntimeConfig { 4654 parse_group_runtime_config_json(&format!( 4655 r#"{{ 4656 "enabled": true, 4657 "canonical_relay_url": "wss://relay.radroots.test", 4658 "relay_secret": "{}", 4659 "owner_pubkeys": ["{}"], 4660 "policy": {{"public_join": true, "invites_enabled": false}} 4661 }}"#, 4662 "7".repeat(64), 4663 owner.as_str() 4664 )) 4665 .expect("groups") 4666 } 4667 4668 fn disabled_groups() -> tangle_groups::GroupRuntimeConfig { 4669 parse_group_runtime_config_json(r#"{"enabled": false}"#).expect("groups") 4670 } 4671 4672 fn signed_auth_event(secret_byte: u8, challenge: &str, created_at: u64) -> Event { 4673 signed_tangle_event_at( 4674 secret_byte, 4675 22_242, 4676 vec![ 4677 Tag::from_parts("relay", &["wss://relay.radroots.test"]).expect("relay"), 4678 Tag::from_parts("challenge", &[challenge]).expect("challenge"), 4679 ], 4680 "", 4681 created_at, 4682 ) 4683 } 4684 4685 fn signed_public_event(secret_byte: u8, kind: u64, tags: Vec<Tag>, content: &str) -> Event { 4686 signed_event_at(secret_byte, kind, tags, content, 1_714_124_433) 4687 } 4688 4689 fn signed_pocket_event( 4690 secret_byte: u8, 4691 kind: u16, 4692 tags: &PocketOwnedTags, 4693 content: &[u8], 4694 ) -> PocketOwnedEvent { 4695 signed_pocket_event_at(secret_byte, kind, tags, content, 1_714_124_433) 4696 } 4697 4698 fn signed_pocket_event_at( 4699 secret_byte: u8, 4700 kind: u16, 4701 tags: &PocketOwnedTags, 4702 content: &[u8], 4703 created_at: u64, 4704 ) -> PocketOwnedEvent { 4705 let secret = format!("{secret_byte:02x}").repeat(32); 4706 RelaySigner::from_secret_hex(&secret) 4707 .expect("signer") 4708 .sign_pocket_event( 4709 PocketKind::from_u16(kind), 4710 tags, 4711 PocketTime::from_u64(created_at), 4712 content, 4713 ) 4714 .expect("pocket event") 4715 } 4716 4717 fn signed_pocket_public_event( 4718 secret_byte: u8, 4719 kind: u32, 4720 tags: Vec<Tag>, 4721 content: &str, 4722 ) -> PocketOwnedEvent { 4723 signed_pocket_event_at_tags(secret_byte, kind, tags, content, 1_714_124_433) 4724 } 4725 4726 fn signed_pocket_group_create_event(secret_byte: u8, group_id: &str) -> PocketOwnedEvent { 4727 signed_pocket_group_create_event_with_tags(secret_byte, group_id, Vec::new(), 1_714_124_433) 4728 } 4729 4730 fn signed_pocket_group_create_event_with_tags( 4731 secret_byte: u8, 4732 group_id: &str, 4733 mut extra_tags: Vec<Tag>, 4734 created_at: u64, 4735 ) -> PocketOwnedEvent { 4736 let mut tags = vec![h(group_id), name(group_id)]; 4737 tags.append(&mut extra_tags); 4738 signed_pocket_event_at_tags(secret_byte, KIND_GROUP_CREATE_GROUP, tags, "", created_at) 4739 } 4740 4741 fn signed_pocket_private_group_create_event( 4742 secret_byte: u8, 4743 group_id: &str, 4744 ) -> PocketOwnedEvent { 4745 signed_pocket_event_at_tags( 4746 secret_byte, 4747 KIND_GROUP_CREATE_GROUP, 4748 vec![h(group_id), name(group_id), private()], 4749 "", 4750 1_714_124_433, 4751 ) 4752 } 4753 4754 fn signed_pocket_event_at_tags( 4755 secret_byte: u8, 4756 kind: u32, 4757 tags: Vec<Tag>, 4758 content: &str, 4759 created_at: u64, 4760 ) -> PocketOwnedEvent { 4761 let tags = pocket_tags_from_protocol(&tags); 4762 signed_pocket_event_at( 4763 secret_byte, 4764 u16::try_from(kind).expect("pocket kind"), 4765 &tags, 4766 content.as_bytes(), 4767 created_at, 4768 ) 4769 } 4770 4771 fn pocket_tags_from_protocol(tags: &[Tag]) -> PocketOwnedTags { 4772 let parts = tags 4773 .iter() 4774 .map(|tag| tag.values().iter().map(String::as_str).collect::<Vec<_>>()) 4775 .collect::<Vec<_>>(); 4776 PocketOwnedTags::new(&parts).expect("pocket tags") 4777 } 4778 4779 fn signed_group_create_event(secret_byte: u8, group_id: &str) -> Event { 4780 signed_group_create_event_with_tags(secret_byte, group_id, Vec::new(), 1_714_124_433) 4781 } 4782 4783 fn signed_group_create_event_with_tags( 4784 secret_byte: u8, 4785 group_id: &str, 4786 mut extra_tags: Vec<Tag>, 4787 created_at: u64, 4788 ) -> Event { 4789 let mut tags = vec![h(group_id), name(group_id)]; 4790 tags.append(&mut extra_tags); 4791 signed_event_at( 4792 secret_byte, 4793 KIND_GROUP_CREATE_GROUP.into(), 4794 tags, 4795 "", 4796 created_at, 4797 ) 4798 } 4799 4800 fn signed_private_group_create_event(secret_byte: u8, group_id: &str) -> Event { 4801 signed_event_at( 4802 secret_byte, 4803 KIND_GROUP_CREATE_GROUP.into(), 4804 vec![h(group_id), name(group_id), private()], 4805 "", 4806 1_714_124_433, 4807 ) 4808 } 4809 4810 fn signed_event_at( 4811 secret_byte: u8, 4812 kind: u64, 4813 tags: Vec<Tag>, 4814 content: &str, 4815 created_at: u64, 4816 ) -> Event { 4817 let pocket = signed_pocket_event_at_tags( 4818 secret_byte, 4819 u32::try_from(kind).expect("kind"), 4820 tags, 4821 content, 4822 created_at, 4823 ); 4824 pocket_event_to_protocol(&pocket) 4825 } 4826 4827 fn signed_tangle_event_at( 4828 secret_byte: u8, 4829 kind: u64, 4830 tags: Vec<Tag>, 4831 content: &str, 4832 created_at: u64, 4833 ) -> Event { 4834 let secret = format!("{secret_byte:02x}").repeat(32); 4835 let signer = RelaySigner::from_secret_hex(&secret).expect("signer"); 4836 let unsigned = UnsignedEvent::new( 4837 signer.public_key().clone(), 4838 UnixTimestamp::new(created_at), 4839 Kind::new(kind).expect("kind"), 4840 tags, 4841 content, 4842 ); 4843 signer.sign_unsigned_event(unsigned) 4844 } 4845 4846 fn pocket_event_id(event: &PocketEvent) -> EventId { 4847 EventId::new(&event.id().as_hex_string()).expect("event id") 4848 } 4849 4850 fn pocket_event_to_protocol(event: &PocketEvent) -> Event { 4851 let tags = event 4852 .tags() 4853 .expect("tags") 4854 .iter() 4855 .map(|tag| { 4856 Tag::new( 4857 tag.map(|value| std::str::from_utf8(value).expect("utf8").to_owned()) 4858 .collect::<Vec<_>>(), 4859 ) 4860 .expect("tag") 4861 }) 4862 .collect::<Vec<_>>(); 4863 Event::new( 4864 pocket_event_id(event), 4865 tangle_protocol::UnsignedEvent::new( 4866 PublicKeyHex::new(&event.pubkey().as_hex_string()).expect("pubkey"), 4867 UnixTimestamp::new(event.created_at().as_u64()), 4868 tangle_protocol::Kind::new(u64::from(event.kind().as_u16())).expect("kind"), 4869 tags, 4870 std::str::from_utf8(event.content()).expect("content"), 4871 ), 4872 SignatureHex::new(&event.sig().to_string()).expect("sig"), 4873 ) 4874 } 4875 4876 fn authenticated_state(secret_byte: u8) -> BaseAuthState { 4877 let mut auth = 4878 BaseAuthState::new("wss://relay.radroots.test", 60, 600).expect("auth state"); 4879 auth.issue_challenge("challenge-a", UnixTimestamp::new(100)) 4880 .expect("challenge"); 4881 let event = signed_auth_event(secret_byte, "challenge-a", 120); 4882 auth.authenticate(&event, UnixTimestamp::new(120)) 4883 .expect("authenticate"); 4884 auth 4885 } 4886 4887 fn count_kind(relay: &BaseRelay, kind: u32) -> u64 { 4888 let subscription_id = SubscriptionId::new(&format!("count-{kind}")).expect("sub"); 4889 let filter = filter_kind(kind); 4890 match relay 4891 .handle_count_protocol(subscription_id, vec![filter]) 4892 .expect("count") 4893 { 4894 RelayMessage::Count { count, .. } => count, 4895 _ => panic!("count response expected"), 4896 } 4897 } 4898 4899 fn count_kind_with_auth(relay: &BaseRelay, kind: u32, auth: &BaseAuthState) -> u64 { 4900 let subscription_id = SubscriptionId::new(&format!("count-auth-{kind}")).expect("sub"); 4901 match relay 4902 .handle_count_with_auth_protocol(subscription_id, vec![filter_kind(kind)], auth) 4903 .expect("count") 4904 { 4905 RelayMessage::Count { count, .. } => count, 4906 _ => panic!("count response expected"), 4907 } 4908 } 4909 4910 fn count_filter(relay: &BaseRelay, subscription_id: &str, filter: Filter) -> u64 { 4911 match relay 4912 .handle_count_protocol( 4913 SubscriptionId::new(subscription_id).expect("sub"), 4914 vec![filter], 4915 ) 4916 .expect("count") 4917 { 4918 RelayMessage::Count { count, .. } => count, 4919 _ => panic!("count response expected"), 4920 } 4921 } 4922 4923 fn count_filter_with_auth( 4924 relay: &BaseRelay, 4925 subscription_id: &str, 4926 filter: Filter, 4927 auth: &BaseAuthState, 4928 ) -> u64 { 4929 match relay 4930 .handle_count_with_auth_protocol( 4931 SubscriptionId::new(subscription_id).expect("sub"), 4932 vec![filter], 4933 auth, 4934 ) 4935 .expect("count") 4936 { 4937 RelayMessage::Count { count, .. } => count, 4938 _ => panic!("count response expected"), 4939 } 4940 } 4941 4942 fn assert_count_without_hll( 4943 relay: &BaseRelay, 4944 subscription_id: &str, 4945 value: serde_json::Value, 4946 auth: Option<&BaseAuthState>, 4947 expected_count: u64, 4948 ) { 4949 let subscription_id = SubscriptionId::new(subscription_id).expect("sub"); 4950 let filter = filter_from_value(&value).expect("filter"); 4951 let message = match auth { 4952 Some(auth) => { 4953 relay.handle_count_with_auth_protocol(subscription_id.clone(), vec![filter], auth) 4954 } 4955 None => relay.handle_count_protocol(subscription_id.clone(), vec![filter]), 4956 } 4957 .expect("count"); 4958 assert_eq!( 4959 message, 4960 RelayMessage::Count { 4961 subscription_id, 4962 count: expected_count, 4963 hll: None 4964 } 4965 ); 4966 } 4967 4968 fn query_filter(relay: &mut BaseRelay, subscription_id: &str, filter: Filter) -> Vec<Event> { 4969 relay 4970 .handle_protocol_req_for_test( 4971 SubscriptionId::new(subscription_id).expect("sub"), 4972 vec![filter], 4973 ) 4974 .expect("query") 4975 .into_iter() 4976 .filter_map(|message| match message { 4977 RelayMessage::Event { event, .. } => Some(event), 4978 _ => None, 4979 }) 4980 .collect() 4981 } 4982 4983 fn filter_kind(kind: u32) -> Filter { 4984 filter_from_value(&serde_json::json!({"kinds":[kind]})).expect("filter") 4985 } 4986 4987 fn filter_group_tag(kind: u32, tag: &str, group_id: &str) -> Filter { 4988 let mut value = serde_json::json!({"kinds":[kind]}); 4989 value 4990 .as_object_mut() 4991 .expect("object") 4992 .insert(format!("#{tag}"), serde_json::json!([group_id])); 4993 filter_from_value(&value).expect("filter") 4994 } 4995 4996 fn pocket_filter_from_value(value: serde_json::Value) -> PocketOwnedFilter { 4997 tangle_filter_to_pocket(&filter_from_value(&value).expect("filter")).expect("pocket") 4998 } 4999 5000 fn hll_target_policy( 5001 relay: &BaseRelay, 5002 value: serde_json::Value, 5003 ) -> BaseRelayCountHllTargetPolicy { 5004 let filter = pocket_filter_from_value(value); 5005 BaseRelay::count_hll_filter_target_policy(relay.groups.as_ref(), &filter) 5006 } 5007 5008 fn count_hll_for_target_policy_test() -> BaseRelayCountHll { 5009 BaseRelayCountHll { 5010 offset: Some(0), 5011 hll: Some(PocketHll8::new()), 5012 suppressed: false, 5013 } 5014 } 5015 5016 fn assert_accepted(message: RelayMessage, event: &Event) { 5017 assert_eq!( 5018 message, 5019 RelayMessage::Ok { 5020 event_id: event.id().clone(), 5021 accepted: true, 5022 message: String::new() 5023 } 5024 ); 5025 } 5026 5027 fn assert_pocket_accepted(message: RelayMessage, event: &PocketEvent) { 5028 assert_eq!( 5029 message, 5030 RelayMessage::Ok { 5031 event_id: pocket_event_id(event), 5032 accepted: true, 5033 message: String::new() 5034 } 5035 ); 5036 } 5037 5038 fn rejected_message(message: RelayMessage) -> String { 5039 match message { 5040 RelayMessage::Ok { 5041 accepted: false, 5042 message, 5043 .. 5044 } => message, 5045 _ => panic!("rejected OK expected"), 5046 } 5047 } 5048 5049 fn assert_member_status( 5050 relay: &BaseRelay, 5051 group_id: &str, 5052 pubkey: &PublicKeyHex, 5053 status: MemberStatus, 5054 ) { 5055 assert_eq!( 5056 relay 5057 .group_projection() 5058 .expect("projection") 5059 .member(&GroupId::new(group_id).expect("group"), pubkey) 5060 .expect("member") 5061 .status(), 5062 status 5063 ); 5064 } 5065 5066 fn has_tag(event: &Event, name: &str, values: &[&str]) -> bool { 5067 event.unsigned().tags().iter().any(|tag| { 5068 tag.values().first().is_some_and(|value| value == name) 5069 && tag.values().len() == values.len() + 1 5070 && values.iter().enumerate().all(|(index, expected)| { 5071 tag.values() 5072 .get(index + 1) 5073 .is_some_and(|value| value == expected) 5074 }) 5075 }) 5076 } 5077 5078 fn h(group_id: &str) -> Tag { 5079 Tag::from_parts("h", &[group_id]).expect("h") 5080 } 5081 5082 fn p(pubkey: &PublicKeyHex) -> Tag { 5083 Tag::from_parts("p", &[pubkey.as_str()]).expect("p") 5084 } 5085 5086 fn e(event_id: &EventId) -> Tag { 5087 Tag::from_parts("e", &[event_id.as_str()]).expect("e") 5088 } 5089 5090 fn name(value: &str) -> Tag { 5091 Tag::from_parts("name", &[value]).expect("name") 5092 } 5093 5094 fn private() -> Tag { 5095 Tag::from_parts("private", &[]).expect("private") 5096 } 5097 5098 fn restricted() -> Tag { 5099 Tag::from_parts("restricted", &[]).expect("restricted") 5100 } 5101 5102 fn hidden() -> Tag { 5103 Tag::from_parts("hidden", &[]).expect("hidden") 5104 } 5105 5106 fn closed() -> Tag { 5107 Tag::from_parts("closed", &[]).expect("closed") 5108 } 5109 5110 fn signer(secret_byte: u8) -> RelaySigner { 5111 RelaySigner::from_secret_hex(&format!("{:02x}", secret_byte).repeat(32)).expect("signer") 5112 } 5113 }