tangle


git clone https://radroots.dev/git/tangle.git
Log | Files | Refs | README | LICENSE

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 }