tangle


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

session.rs (84795B)


      1 #![forbid(unsafe_code)]
      2 
      3 use crate::{
      4     client_message::{RuntimeClientMessage, parse_runtime_client_message},
      5     errors::BaseRelayError,
      6     event_bus::{TangleEventReceiveError, TangleEventReceiver},
      7     logging,
      8     relay::{
      9         auth::{BaseAuthState, generate_auth_challenge},
     10         core::BaseRelay,
     11         live::{CloseResult, LiveSubscriptionSet},
     12         outbound::RuntimeRelayMessage,
     13     },
     14     resource_limits::{RelayResourceLimiter, RelaySubscriptionPermit},
     15     runtime::{
     16         RelayProjectionContext, RelayRuntimeHandle, TangleClientMessageMetricKind,
     17         TangleClientRateLimitContext, TangleRuntimeLimits,
     18     },
     19 };
     20 use axum::extract::ws::{CloseFrame, Message, Utf8Bytes, WebSocket};
     21 use std::{
     22     collections::BTreeMap,
     23     future::pending,
     24     net::IpAddr,
     25     sync::atomic::{AtomicU64, Ordering},
     26     time::{Duration, Instant, SystemTime, UNIX_EPOCH},
     27 };
     28 use tangle_protocol::{RelayMessage, SubscriptionId, UnixTimestamp};
     29 use tangle_store_pocket::PocketOwnedFilter;
     30 use tokio::{
     31     sync::{mpsc, watch},
     32     time::{MissedTickBehavior, interval_at, timeout},
     33 };
     34 
     35 #[derive(Debug)]
     36 pub struct TangleWebSocketSession {
     37     connection_id: u64,
     38     peer_ip: Option<IpAddr>,
     39     connected_at: Instant,
     40     outbound: TangleOutboundSender,
     41     outbound_receiver: mpsc::Receiver<Message>,
     42     shutdown: watch::Receiver<bool>,
     43     runtime: RelayRuntimeHandle,
     44     limits: TangleRuntimeLimits,
     45     auth: BaseAuthState,
     46     subscriptions: LiveSubscriptionSet,
     47     resource_limiter: Option<RelayResourceLimiter>,
     48     subscription_permits: BTreeMap<SubscriptionId, RelaySubscriptionPermit>,
     49     events: TangleEventReceiver,
     50     projection: RelayProjectionContext,
     51     keepalive: Option<TangleWebSocketKeepaliveConfig>,
     52 }
     53 
     54 static NEXT_TANGLE_CONNECTION_ID: AtomicU64 = AtomicU64::new(1);
     55 
     56 #[derive(Debug, Clone, Default)]
     57 pub struct TangleWebSocketSessionOptions {
     58     peer_ip: Option<IpAddr>,
     59     resource_limiter: Option<RelayResourceLimiter>,
     60     projection: RelayProjectionContext,
     61     keepalive: Option<TangleWebSocketKeepaliveConfig>,
     62 }
     63 
     64 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
     65 pub struct TangleWebSocketKeepaliveConfig {
     66     ping_interval: Duration,
     67     peer_timeout: Duration,
     68     write_timeout: Duration,
     69 }
     70 
     71 impl TangleWebSocketKeepaliveConfig {
     72     pub fn new(
     73         ping_interval: Duration,
     74         peer_timeout: Duration,
     75         write_timeout: Duration,
     76     ) -> Result<Self, BaseRelayError> {
     77         if ping_interval.is_zero() || peer_timeout.is_zero() || write_timeout.is_zero() {
     78             return Err(BaseRelayError::invalid(
     79                 "websocket keepalive durations must be greater than zero",
     80             ));
     81         }
     82         if ping_interval >= peer_timeout {
     83             return Err(BaseRelayError::invalid(
     84                 "websocket ping interval must be shorter than peer timeout",
     85             ));
     86         }
     87         if write_timeout > peer_timeout {
     88             return Err(BaseRelayError::invalid(
     89                 "websocket write timeout must not exceed peer timeout",
     90             ));
     91         }
     92         Ok(Self {
     93             ping_interval,
     94             peer_timeout,
     95             write_timeout,
     96         })
     97     }
     98 
     99     pub fn ping_interval(self) -> Duration {
    100         self.ping_interval
    101     }
    102 
    103     pub fn peer_timeout(self) -> Duration {
    104         self.peer_timeout
    105     }
    106 
    107     pub fn write_timeout(self) -> Duration {
    108         self.write_timeout
    109     }
    110 }
    111 
    112 impl TangleWebSocketSessionOptions {
    113     pub fn new() -> Self {
    114         Self::default()
    115     }
    116 
    117     pub fn with_peer_ip(mut self, peer_ip: Option<IpAddr>) -> Self {
    118         self.peer_ip = peer_ip;
    119         self
    120     }
    121 
    122     pub fn with_resource_limiter(mut self, resource_limiter: Option<RelayResourceLimiter>) -> Self {
    123         self.resource_limiter = resource_limiter;
    124         self
    125     }
    126 
    127     pub fn with_projection_context(mut self, projection: RelayProjectionContext) -> Self {
    128         self.projection = projection;
    129         self
    130     }
    131 
    132     pub fn with_keepalive(mut self, keepalive: Option<TangleWebSocketKeepaliveConfig>) -> Self {
    133         self.keepalive = keepalive;
    134         self
    135     }
    136 }
    137 
    138 impl TangleWebSocketSession {
    139     pub fn new(
    140         limits: TangleRuntimeLimits,
    141         shutdown: watch::Receiver<bool>,
    142         runtime: RelayRuntimeHandle,
    143         auth: BaseAuthState,
    144         events: TangleEventReceiver,
    145     ) -> Result<Self, BaseRelayError> {
    146         Self::new_with_peer(limits, shutdown, runtime, auth, events, None)
    147     }
    148 
    149     pub fn new_with_peer(
    150         limits: TangleRuntimeLimits,
    151         shutdown: watch::Receiver<bool>,
    152         runtime: RelayRuntimeHandle,
    153         auth: BaseAuthState,
    154         events: TangleEventReceiver,
    155         peer_ip: Option<IpAddr>,
    156     ) -> Result<Self, BaseRelayError> {
    157         Self::new_with_peer_and_resources(limits, shutdown, runtime, auth, events, peer_ip, None)
    158     }
    159 
    160     pub fn new_with_peer_and_resources(
    161         limits: TangleRuntimeLimits,
    162         shutdown: watch::Receiver<bool>,
    163         runtime: RelayRuntimeHandle,
    164         auth: BaseAuthState,
    165         events: TangleEventReceiver,
    166         peer_ip: Option<IpAddr>,
    167         resource_limiter: Option<RelayResourceLimiter>,
    168     ) -> Result<Self, BaseRelayError> {
    169         Self::new_with_options(
    170             limits,
    171             shutdown,
    172             runtime,
    173             auth,
    174             events,
    175             TangleWebSocketSessionOptions::new()
    176                 .with_peer_ip(peer_ip)
    177                 .with_resource_limiter(resource_limiter),
    178         )
    179     }
    180 
    181     pub fn new_with_options(
    182         limits: TangleRuntimeLimits,
    183         shutdown: watch::Receiver<bool>,
    184         runtime: RelayRuntimeHandle,
    185         auth: BaseAuthState,
    186         events: TangleEventReceiver,
    187         options: TangleWebSocketSessionOptions,
    188     ) -> Result<Self, BaseRelayError> {
    189         let outbound_queue_capacity = limits.outbound_queue_capacity();
    190         let (sender, receiver) = mpsc::channel(outbound_queue_capacity);
    191         let subscriptions = LiveSubscriptionSet::new(
    192             limits.base_relay_limits().max_pending_events(),
    193             limits.base_relay_limits().max_subscriptions(),
    194         )?;
    195         Ok(Self {
    196             connection_id: NEXT_TANGLE_CONNECTION_ID.fetch_add(1, Ordering::Relaxed),
    197             peer_ip: options.peer_ip,
    198             connected_at: Instant::now(),
    199             outbound: TangleOutboundSender {
    200                 sender,
    201                 capacity: outbound_queue_capacity,
    202             },
    203             outbound_receiver: receiver,
    204             shutdown,
    205             runtime,
    206             limits,
    207             auth,
    208             subscriptions,
    209             resource_limiter: options.resource_limiter,
    210             subscription_permits: BTreeMap::new(),
    211             events,
    212             projection: options.projection,
    213             keepalive: options.keepalive,
    214         })
    215     }
    216 
    217     pub fn connected_at(&self) -> Instant {
    218         self.connected_at
    219     }
    220 
    221     pub fn outbound(&self) -> TangleOutboundSender {
    222         self.outbound.clone()
    223     }
    224 
    225     pub fn shutdown_requested(&self) -> bool {
    226         *self.shutdown.borrow()
    227     }
    228 
    229     #[cfg(test)]
    230     fn active_subscription_count(&self) -> usize {
    231         self.subscriptions.active_count()
    232     }
    233 
    234     pub async fn run(mut self, mut socket: WebSocket) {
    235         let metrics = self.runtime.metrics();
    236         metrics.record_session_opened();
    237         logging::log_websocket_session_opened(self.connection_id, self.peer_ip);
    238         if !self.issue_auth_challenge() {
    239             let closed_subscriptions = self.close_all_subscriptions();
    240             metrics.record_subscriptions_closed(closed_subscriptions);
    241             metrics.record_session_closed();
    242             metrics.record_event_bus_receivers(
    243                 metrics.event_bus_receivers_current().saturating_sub(1),
    244             );
    245             logging::log_websocket_session_closed(
    246                 self.connection_id,
    247                 self.peer_ip,
    248                 closed_subscriptions,
    249             );
    250             return;
    251         }
    252         let mut last_peer_activity = Instant::now();
    253         let mut keepalive_interval = self.keepalive.map(|config| {
    254             let mut interval = interval_at(
    255                 tokio::time::Instant::now() + config.ping_interval(),
    256                 config.ping_interval(),
    257             );
    258             interval.set_missed_tick_behavior(MissedTickBehavior::Delay);
    259             interval
    260         });
    261         loop {
    262             if self.shutdown_requested() {
    263                 let _ = self
    264                     .send_socket_message(&mut socket, Message::Close(None))
    265                     .await;
    266                 break;
    267             }
    268             tokio::select! {
    269                 incoming = socket.recv() => {
    270                     match incoming {
    271                         Some(Ok(Message::Close(_))) | Some(Err(_)) | None => break,
    272                         Some(Ok(Message::Ping(payload))) => {
    273                             last_peer_activity = Instant::now();
    274                             if !self.send_socket_message(&mut socket, Message::Pong(payload)).await {
    275                                 break;
    276                             }
    277                         }
    278                         Some(Ok(Message::Pong(_))) => {
    279                             last_peer_activity = Instant::now();
    280                         }
    281                         Some(Ok(message)) => {
    282                             last_peer_activity = Instant::now();
    283                             match self.handle_incoming_message(message).await {
    284                                 TangleSessionControl::Continue => {}
    285                                 TangleSessionControl::Close(message) => {
    286                                     let _ = self.send_socket_message(&mut socket, message).await;
    287                                     break;
    288                                 }
    289                                 TangleSessionControl::Stop => break,
    290                             }
    291                         }
    292                     }
    293                 }
    294                 outgoing = self.outbound_receiver.recv() => {
    295                     let Some(message) = outgoing else {
    296                         break;
    297                     };
    298                     if !self.send_socket_message(&mut socket, message).await {
    299                         break;
    300                     }
    301                 }
    302                 event = self.events.recv() => {
    303                     match self.handle_event_receive_result(event).await {
    304                         TangleSessionControl::Continue => {}
    305                         TangleSessionControl::Close(message) => {
    306                             let _ = self.send_socket_message(&mut socket, message).await;
    307                             break;
    308                         }
    309                         TangleSessionControl::Stop => break,
    310                     }
    311                 }
    312                 changed = self.shutdown.changed() => {
    313                     if changed.is_err() || self.shutdown_requested() {
    314                         let _ = self.send_socket_message(&mut socket, Message::Close(None)).await;
    315                         break;
    316                     }
    317                 }
    318                 _ = keepalive_tick(&mut keepalive_interval) => {
    319                     let Some(config) = self.keepalive else {
    320                         continue;
    321                     };
    322                     if last_peer_activity.elapsed() >= config.peer_timeout() {
    323                         let _ = self
    324                             .send_socket_message(&mut socket, keepalive_timeout_close_message())
    325                             .await;
    326                         break;
    327                     }
    328                     if !self
    329                         .send_socket_message(&mut socket, Message::Ping(Vec::new().into()))
    330                         .await
    331                     {
    332                         break;
    333                     }
    334                 }
    335             }
    336         }
    337         let closed_subscriptions = self.close_all_subscriptions();
    338         metrics.record_subscriptions_closed(closed_subscriptions);
    339         metrics.record_session_closed();
    340         metrics.record_event_bus_receivers(metrics.event_bus_receivers_current().saturating_sub(1));
    341         logging::log_websocket_session_closed(
    342             self.connection_id,
    343             self.peer_ip,
    344             closed_subscriptions,
    345         );
    346     }
    347 
    348     async fn handle_event_receive_result(
    349         &mut self,
    350         result: Result<tangle_groups::StoreOffset, TangleEventReceiveError>,
    351     ) -> TangleSessionControl {
    352         match result {
    353             Ok(offset) => self.handle_event_offset(offset).await,
    354             Err(TangleEventReceiveError::Lagged(skipped)) => {
    355                 self.runtime.metrics().record_event_bus_lagged(skipped);
    356                 TangleSessionControl::Close(event_stream_lag_close_message())
    357             }
    358             Err(TangleEventReceiveError::Closed) => TangleSessionControl::Stop,
    359             Err(TangleEventReceiveError::Empty) => TangleSessionControl::Continue,
    360         }
    361     }
    362 
    363     async fn handle_event_offset(
    364         &mut self,
    365         offset: tangle_groups::StoreOffset,
    366     ) -> TangleSessionControl {
    367         let runtime = self.runtime.clone();
    368         let auth = self.auth.clone();
    369         let replies = match runtime
    370             .fanout_event_offset_with_projection_context(
    371                 offset,
    372                 &mut self.subscriptions,
    373                 &auth,
    374                 &self.projection,
    375             )
    376             .await
    377         {
    378             Ok(replies) => replies,
    379             Err(error) => vec![RelayMessage::Notice(error.prefixed_message()).into()],
    380         };
    381         for reply in replies {
    382             if let Err(control) = self.enqueue_relay_message(reply) {
    383                 return control;
    384             }
    385         }
    386         TangleSessionControl::Continue
    387     }
    388 
    389     async fn handle_incoming_message(&mut self, message: Message) -> TangleSessionControl {
    390         match message {
    391             Message::Text(raw) => self.dispatch_text(raw.as_str()).await,
    392             Message::Binary(_) => self
    393                 .enqueue_relay_message(
    394                     RelayMessage::Notice("invalid: client message must be a text frame".to_owned())
    395                         .into(),
    396                 )
    397                 .map(|_| TangleSessionControl::Continue)
    398                 .unwrap_or_else(|control| control),
    399             Message::Ping(_) | Message::Pong(_) => TangleSessionControl::Continue,
    400             Message::Close(_) => TangleSessionControl::Stop,
    401         }
    402     }
    403 
    404     fn issue_auth_challenge(&mut self) -> bool {
    405         let message = generate_auth_challenge()
    406             .and_then(|challenge| {
    407                 self.auth
    408                     .issue_challenge(challenge, current_unix_timestamp())
    409             })
    410             .unwrap_or_else(|error| RelayMessage::Notice(error.prefixed_message()));
    411         self.send_relay_message(message.into()).is_ok()
    412     }
    413 
    414     async fn dispatch_text(&mut self, raw: &str) -> TangleSessionControl {
    415         if raw.len() > self.limits.max_message_length() {
    416             return self
    417                 .enqueue_relay_message(
    418                     RelayMessage::Notice(format!(
    419                         "invalid: client message length exceeds runtime max_message_length {}",
    420                         self.limits.max_message_length()
    421                     ))
    422                     .into(),
    423                 )
    424                 .map(|_| TangleSessionControl::Continue)
    425                 .unwrap_or_else(|control| control);
    426         }
    427         let replies = match parse_runtime_client_message(raw) {
    428             Ok(message) => match self.handle_client_message(message).await {
    429                 Ok(replies) => replies,
    430                 Err(error) => vec![RelayMessage::Notice(error.prefixed_message()).into()],
    431             },
    432             Err(error) => vec![RelayMessage::Notice(format!("invalid: {error}")).into()],
    433         };
    434         for reply in replies {
    435             if let Err(control) = self.enqueue_relay_message(reply) {
    436                 return control;
    437             }
    438         }
    439         TangleSessionControl::Continue
    440     }
    441 
    442     async fn handle_client_message(
    443         &mut self,
    444         message: RuntimeClientMessage,
    445     ) -> Result<Vec<RuntimeRelayMessage>, BaseRelayError> {
    446         match message {
    447             RuntimeClientMessage::Req {
    448                 subscription_id,
    449                 filters,
    450                 search_present,
    451             } => {
    452                 self.handle_req(subscription_id, filters, search_present)
    453                     .await
    454             }
    455             RuntimeClientMessage::Count {
    456                 subscription_id,
    457                 filters,
    458                 search_present,
    459             } => {
    460                 let context = self.client_rate_limit_context();
    461                 self.runtime
    462                     .handle_client_message_with_rate_limit_context(
    463                         RuntimeClientMessage::Count {
    464                             subscription_id,
    465                             filters,
    466                             search_present,
    467                         },
    468                         &mut self.auth,
    469                         context,
    470                         current_unix_timestamp(),
    471                     )
    472                     .await
    473             }
    474             RuntimeClientMessage::Close(subscription_id) => {
    475                 let metrics = self.runtime.metrics();
    476                 metrics.record_client_message(TangleClientMessageMetricKind::Close);
    477                 self.limits
    478                     .base_relay_limits()
    479                     .validate_subscription_id(&subscription_id)?;
    480                 if self.subscriptions.close(&subscription_id) == CloseResult::Closed {
    481                     self.subscription_permits.remove(&subscription_id);
    482                     metrics.record_subscriptions_closed(1);
    483                 }
    484                 Ok(Vec::new())
    485             }
    486             message => {
    487                 let context = self.client_rate_limit_context();
    488                 self.runtime
    489                     .handle_client_message_with_rate_limit_context(
    490                         message,
    491                         &mut self.auth,
    492                         context,
    493                         current_unix_timestamp(),
    494                     )
    495                     .await
    496             }
    497         }
    498     }
    499 
    500     async fn handle_req(
    501         &mut self,
    502         subscription_id: SubscriptionId,
    503         filters: Vec<PocketOwnedFilter>,
    504         search_present: bool,
    505     ) -> Result<Vec<RuntimeRelayMessage>, BaseRelayError> {
    506         let metrics = self.runtime.metrics();
    507         metrics.record_client_message(TangleClientMessageMetricKind::Req);
    508         self.limits
    509             .base_relay_limits()
    510             .validate_subscription_id(&subscription_id)?;
    511         self.limits
    512             .base_relay_limits()
    513             .validate_pocket_filters(&filters)?;
    514         if let Some(message) =
    515             BaseRelay::unsupported_search_present_closed(&subscription_id, search_present)
    516         {
    517             return Ok(vec![message.into()]);
    518         }
    519         if let Some(message) = self
    520             .runtime
    521             .rate_limit_req_pocket(
    522                 &subscription_id,
    523                 &filters,
    524                 &self.auth,
    525                 self.client_rate_limit_context(),
    526                 current_unix_timestamp(),
    527             )
    528             .await
    529         {
    530             return Ok(vec![message.into()]);
    531         }
    532         let should_subscribe = !pocket_filters_are_complete(&filters);
    533         let already_subscribed = self.subscriptions.contains(&subscription_id);
    534         if should_subscribe {
    535             self.subscriptions
    536                 .ensure_can_subscribe(&subscription_id, &filters)?;
    537         }
    538         let report = self
    539             .runtime
    540             .query_req_with_auth_report_with_projection_context(
    541                 subscription_id.clone(),
    542                 filters.clone(),
    543                 search_present,
    544                 &self.auth,
    545                 &self.projection,
    546             )
    547             .await?;
    548         let closes_subscription = report.group_read_denied();
    549         let replies = report.into_messages();
    550         if should_subscribe && !closes_subscription {
    551             let host_permit = if already_subscribed {
    552                 None
    553             } else {
    554                 self.resource_limiter
    555                     .as_ref()
    556                     .map(|resources| resources.try_open_subscriptions(1))
    557                     .transpose()?
    558             };
    559             self.subscriptions
    560                 .subscribe(subscription_id.clone(), filters)?;
    561             if let Some(permit) = host_permit {
    562                 self.subscription_permits
    563                     .insert(subscription_id.clone(), permit);
    564             }
    565             if !already_subscribed {
    566                 metrics.record_subscription_opened();
    567                 logging::log_subscription_opened(self.connection_id, &subscription_id);
    568             }
    569         }
    570         Ok(replies)
    571     }
    572 
    573     fn close_all_subscriptions(&mut self) -> usize {
    574         let closed = self.subscriptions.close_all();
    575         self.subscription_permits.clear();
    576         closed
    577     }
    578 
    579     fn client_rate_limit_context(&self) -> TangleClientRateLimitContext {
    580         TangleClientRateLimitContext::new(self.peer_ip, Some(self.connection_id))
    581     }
    582 
    583     fn send_relay_message(&self, message: RuntimeRelayMessage) -> Result<(), TangleSessionControl> {
    584         let text = self
    585             .runtime
    586             .sanitize_public_message(message)
    587             .encode()
    588             .map_err(|_| TangleSessionControl::Close(outbound_encode_close_message()))?;
    589         self.outbound
    590             .try_send(Message::Text(text.into()))
    591             .map_err(|error| self.outbound_queue_error_control(error))
    592     }
    593 
    594     fn enqueue_relay_message(
    595         &self,
    596         message: RuntimeRelayMessage,
    597     ) -> Result<(), TangleSessionControl> {
    598         self.send_relay_message(message)
    599     }
    600 
    601     fn outbound_queue_error_control(
    602         &self,
    603         error: TangleOutboundQueueError,
    604     ) -> TangleSessionControl {
    605         match error {
    606             TangleOutboundQueueError::Full => {
    607                 self.runtime.metrics().record_outbound_queue_full_close();
    608                 TangleSessionControl::Close(outbound_queue_full_close_message())
    609             }
    610             TangleOutboundQueueError::Closed => TangleSessionControl::Stop,
    611         }
    612     }
    613 
    614     async fn send_socket_message(&self, socket: &mut WebSocket, message: Message) -> bool {
    615         match self.keepalive {
    616             Some(config) => timeout(config.write_timeout(), socket.send(message))
    617                 .await
    618                 .is_ok_and(|result| result.is_ok()),
    619             None => socket.send(message).await.is_ok(),
    620         }
    621     }
    622 }
    623 
    624 async fn keepalive_tick(interval: &mut Option<tokio::time::Interval>) {
    625     match interval {
    626         Some(interval) => {
    627             interval.tick().await;
    628         }
    629         None => pending::<()>().await,
    630     }
    631 }
    632 
    633 #[derive(Debug, Clone, PartialEq, Eq)]
    634 enum TangleSessionControl {
    635     Continue,
    636     Close(Message),
    637     Stop,
    638 }
    639 
    640 fn event_stream_lag_close_message() -> Message {
    641     Message::Close(Some(CloseFrame {
    642         code: 1008,
    643         reason: Utf8Bytes::from_static("event stream lagged; reconnect required"),
    644     }))
    645 }
    646 
    647 fn outbound_queue_full_close_message() -> Message {
    648     Message::Close(Some(CloseFrame {
    649         code: 1013,
    650         reason: Utf8Bytes::from_static("outbound queue full; reconnect required"),
    651     }))
    652 }
    653 
    654 fn outbound_encode_close_message() -> Message {
    655     Message::Close(Some(CloseFrame {
    656         code: 1011,
    657         reason: Utf8Bytes::from_static("outbound relay message encode failed"),
    658     }))
    659 }
    660 
    661 fn keepalive_timeout_close_message() -> Message {
    662     Message::Close(Some(CloseFrame {
    663         code: 1001,
    664         reason: Utf8Bytes::from_static("websocket peer timeout"),
    665     }))
    666 }
    667 
    668 #[derive(Debug, Clone)]
    669 pub struct TangleOutboundSender {
    670     sender: mpsc::Sender<Message>,
    671     capacity: usize,
    672 }
    673 
    674 impl TangleOutboundSender {
    675     pub fn capacity(&self) -> usize {
    676         self.capacity
    677     }
    678 
    679     pub fn try_send(&self, message: Message) -> Result<(), TangleOutboundQueueError> {
    680         self.sender.try_send(message).map_err(Into::into)
    681     }
    682 }
    683 
    684 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
    685 pub enum TangleOutboundQueueError {
    686     Full,
    687     Closed,
    688 }
    689 
    690 impl From<mpsc::error::TrySendError<Message>> for TangleOutboundQueueError {
    691     fn from(error: mpsc::error::TrySendError<Message>) -> Self {
    692         match error {
    693             mpsc::error::TrySendError::Full(_) => Self::Full,
    694             mpsc::error::TrySendError::Closed(_) => Self::Closed,
    695         }
    696     }
    697 }
    698 
    699 fn current_unix_timestamp() -> UnixTimestamp {
    700     UnixTimestamp::new(
    701         SystemTime::now()
    702             .duration_since(UNIX_EPOCH)
    703             .map(|duration| duration.as_secs())
    704             .unwrap_or(0),
    705     )
    706 }
    707 
    708 fn pocket_filters_are_complete(filters: &[PocketOwnedFilter]) -> bool {
    709     !filters.is_empty() && filters.iter().all(|filter| filter.completes())
    710 }
    711 
    712 #[cfg(test)]
    713 impl TangleWebSocketSession {
    714     async fn handle_protocol_client_message_for_test(
    715         &mut self,
    716         message: tangle_protocol::ClientMessage,
    717     ) -> Result<Vec<RelayMessage>, BaseRelayError> {
    718         let messages = self
    719             .handle_client_message(protocol_client_message_to_runtime_for_session_test(
    720                 message,
    721             )?)
    722             .await?;
    723         crate::relay::outbound::protocol_messages_for_test(messages)
    724     }
    725 }
    726 
    727 #[cfg(test)]
    728 fn protocol_client_message_to_runtime_for_session_test(
    729     message: tangle_protocol::ClientMessage,
    730 ) -> Result<RuntimeClientMessage, BaseRelayError> {
    731     match message {
    732         tangle_protocol::ClientMessage::Event(event) => Ok(RuntimeClientMessage::Event(
    733             crate::pocket_conversion::tangle_event_to_pocket(&event)?,
    734         )),
    735         tangle_protocol::ClientMessage::Req {
    736             subscription_id,
    737             filters,
    738         } => Ok(RuntimeClientMessage::Req {
    739             subscription_id,
    740             search_present: filters.iter().any(|filter| filter.search().is_some()),
    741             filters: filters
    742                 .iter()
    743                 .map(crate::pocket_conversion::tangle_filter_to_pocket)
    744                 .collect::<Result<Vec<_>, _>>()?,
    745         }),
    746         tangle_protocol::ClientMessage::Count {
    747             subscription_id,
    748             filters,
    749         } => Ok(RuntimeClientMessage::Count {
    750             subscription_id,
    751             search_present: filters.iter().any(|filter| filter.search().is_some()),
    752             filters: filters
    753                 .iter()
    754                 .map(crate::pocket_conversion::tangle_filter_to_pocket)
    755                 .collect::<Result<Vec<_>, _>>()?,
    756         }),
    757         tangle_protocol::ClientMessage::Close(subscription_id) => {
    758             Ok(RuntimeClientMessage::Close(subscription_id))
    759         }
    760         tangle_protocol::ClientMessage::Auth(event) => Ok(RuntimeClientMessage::Auth(
    761             crate::pocket_conversion::tangle_event_to_pocket(&event)?,
    762         )),
    763         tangle_protocol::ClientMessage::NegOpen {
    764             subscription_id,
    765             filter,
    766             message,
    767         } => Ok(RuntimeClientMessage::NegOpen {
    768             subscription_id,
    769             filter: crate::pocket_conversion::tangle_filter_to_pocket(&filter)?,
    770             message,
    771         }),
    772         tangle_protocol::ClientMessage::NegMsg {
    773             subscription_id,
    774             message,
    775         } => Ok(RuntimeClientMessage::NegMsg {
    776             subscription_id,
    777             message,
    778         }),
    779         tangle_protocol::ClientMessage::NegClose(subscription_id) => {
    780             Ok(RuntimeClientMessage::NegClose(subscription_id))
    781         }
    782     }
    783 }
    784 
    785 #[cfg(test)]
    786 mod tests {
    787     use super::{
    788         TangleOutboundQueueError, TangleSessionControl, TangleWebSocketKeepaliveConfig,
    789         TangleWebSocketSession, current_unix_timestamp, event_stream_lag_close_message,
    790         keepalive_timeout_close_message, outbound_queue_full_close_message,
    791     };
    792     use crate::{
    793         config::{BaseRelayRuntimeConfig, parse_base_relay_runtime_config_json},
    794         errors::BaseRelayError,
    795         event_bus::TangleEventReceiver,
    796         rate_limits::{TangleRateLimitKey, TangleRateLimitScope},
    797         relay::core::{BaseRelayLimitSettings, BaseRelayLimits},
    798         runtime::{RelayRuntime, RelayRuntimeHandle, TangleRuntimeLimits, TangleShutdownSignal},
    799     };
    800     use axum::extract::ws::Message;
    801     use serde_json::json;
    802     use std::{
    803         path::{Path, PathBuf},
    804         time::Duration,
    805     };
    806     use tangle_crypto::RelaySigner;
    807     use tangle_groups::{KIND_GROUP_CREATE_GROUP, StoreOffset};
    808     use tangle_protocol::{
    809         ClientMessage, Event, EventId, Filter, Kind, PublicKeyHex, RelayMessage, SignatureHex,
    810         SubscriptionId, Tag, UnixTimestamp, UnsignedEvent, event_to_value, filter_from_value,
    811     };
    812     use tangle_store_pocket::{
    813         PocketEvent, PocketKind, PocketOwnedEvent, PocketOwnedTags, PocketTime,
    814     };
    815     use tangle_test_support::FixtureKey;
    816 
    817     #[test]
    818     fn websocket_session_records_connection_time() {
    819         let before = std::time::Instant::now();
    820         let shutdown = TangleShutdownSignal::new();
    821         let (runtime, auth, events) = session_runtime("records-connection-time");
    822         let session = TangleWebSocketSession::new(
    823             session_limits(8),
    824             shutdown.subscribe(),
    825             runtime,
    826             auth,
    827             events,
    828         )
    829         .expect("session");
    830 
    831         assert!(session.connected_at() >= before);
    832     }
    833 
    834     #[test]
    835     fn websocket_keepalive_config_is_bounded_and_explicit() {
    836         let config = TangleWebSocketKeepaliveConfig::new(
    837             Duration::from_secs(30),
    838             Duration::from_secs(90),
    839             Duration::from_secs(5),
    840         )
    841         .expect("keepalive");
    842         assert_eq!(config.ping_interval(), Duration::from_secs(30));
    843         assert_eq!(config.peer_timeout(), Duration::from_secs(90));
    844         assert_eq!(config.write_timeout(), Duration::from_secs(5));
    845         assert!(
    846             TangleWebSocketKeepaliveConfig::new(
    847                 Duration::ZERO,
    848                 Duration::from_secs(90),
    849                 Duration::from_secs(5),
    850             )
    851             .is_err()
    852         );
    853         assert!(
    854             TangleWebSocketKeepaliveConfig::new(
    855                 Duration::from_secs(90),
    856                 Duration::from_secs(90),
    857                 Duration::from_secs(5),
    858             )
    859             .is_err()
    860         );
    861         assert!(
    862             TangleWebSocketKeepaliveConfig::new(
    863                 Duration::from_secs(30),
    864                 Duration::from_secs(90),
    865                 Duration::from_secs(91),
    866             )
    867             .is_err()
    868         );
    869         let Message::Close(Some(frame)) = keepalive_timeout_close_message() else {
    870             panic!("keepalive close frame")
    871         };
    872         assert_eq!(frame.code, 1001);
    873         assert_eq!(frame.reason, "websocket peer timeout");
    874     }
    875 
    876     #[test]
    877     fn websocket_session_limit_config_rejects_zero_outbound_capacity() {
    878         assert!(session_limits_result(0).is_err());
    879     }
    880 
    881     #[test]
    882     fn websocket_session_observes_shutdown_request() {
    883         let shutdown = TangleShutdownSignal::new();
    884         let (runtime, auth, events) = session_runtime("observes-shutdown");
    885         let session = TangleWebSocketSession::new(
    886             session_limits(8),
    887             shutdown.subscribe(),
    888             runtime,
    889             auth,
    890             events,
    891         )
    892         .expect("session");
    893 
    894         assert!(!session.shutdown_requested());
    895 
    896         shutdown.request_shutdown();
    897 
    898         assert!(session.shutdown_requested());
    899     }
    900 
    901     #[tokio::test]
    902     async fn websocket_session_rejects_overlong_text_before_parsing() {
    903         let shutdown = TangleShutdownSignal::new();
    904         let (runtime, auth, events) = session_runtime("overlong-text");
    905         let mut session = TangleWebSocketSession::new(
    906             session_limits_with_message_length(8, 8),
    907             shutdown.subscribe(),
    908             runtime,
    909             auth,
    910             events,
    911         )
    912         .expect("session");
    913 
    914         assert_eq!(
    915             session.dispatch_text("123456789").await,
    916             TangleSessionControl::Continue
    917         );
    918         let message = session.outbound_receiver.try_recv().expect("notice");
    919         let Message::Text(text) = message else {
    920             panic!("expected text notice")
    921         };
    922         assert_eq!(
    923             text.as_str(),
    924             "[\"NOTICE\",\"invalid: client message length exceeds runtime max_message_length 8\"]"
    925         );
    926     }
    927 
    928     #[tokio::test]
    929     async fn websocket_session_preserves_chorus_malformed_message_parity() {
    930         let shutdown = TangleShutdownSignal::new();
    931         let (runtime, auth, events) = session_runtime("chorus-malformed-parity");
    932         let mut session = TangleWebSocketSession::new(
    933             session_limits(16),
    934             shutdown.subscribe(),
    935             runtime,
    936             auth,
    937             events,
    938         )
    939         .expect("session");
    940         let event = tangle_v2_event(FixtureKey::Member, 1_714_124_433, 1, Vec::new(), "parity")
    941             .expect("event");
    942         for (raw, expected) in [
    943             ("{", None),
    944             (
    945                 "[\"NOTICE\",\"client\"]",
    946                 Some("[\"NOTICE\",\"invalid: client message command `NOTICE` is unsupported\"]"),
    947             ),
    948             (
    949                 "[\"NEG-OPEN\",\"sub\",{}]",
    950                 Some(
    951                     "[\"NOTICE\",\"invalid: NEG-OPEN client message must contain a subscription id, filter, and message\"]",
    952                 ),
    953             ),
    954             (
    955                 "[\"REQ\"]",
    956                 Some(
    957                     "[\"NOTICE\",\"invalid: REQ client message must contain a subscription id and filters\"]",
    958                 ),
    959             ),
    960             (
    961                 "[\"CLOSE\",1]",
    962                 Some("[\"NOTICE\",\"invalid: CLOSE subscription id must be a string\"]"),
    963             ),
    964         ] {
    965             assert_eq!(
    966                 session.dispatch_text(raw).await,
    967                 TangleSessionControl::Continue
    968             );
    969             let text = take_outbound_text(&mut session);
    970             if let Some(expected) = expected {
    971                 assert_eq!(text, expected);
    972             } else {
    973                 assert!(text.starts_with("[\"NOTICE\",\"invalid: client message JSON is invalid:"));
    974             }
    975         }
    976 
    977         assert_eq!(
    978             session
    979                 .dispatch_text("[\"REQ\",\"sub-search\",{\"search\":\"carrots\"}]")
    980                 .await,
    981             TangleSessionControl::Continue
    982         );
    983         assert_eq!(
    984             take_outbound_text(&mut session),
    985             "[\"CLOSED\",\"sub-search\",\"unsupported: search filters are not supported\"]"
    986         );
    987 
    988         assert_eq!(
    989             session
    990                 .dispatch_text(&json!(["EVENT", event_to_value(&event)]).to_string())
    991                 .await,
    992             TangleSessionControl::Continue
    993         );
    994         assert_eq!(
    995             take_outbound_text(&mut session),
    996             format!("[\"OK\",\"{}\",true,\"\"]", event.id().as_str())
    997         );
    998     }
    999 
   1000     #[tokio::test]
   1001     async fn websocket_session_returns_disabled_negentropy_errors() {
   1002         let shutdown = TangleShutdownSignal::new();
   1003         let (runtime, auth, events) = session_runtime("disabled-negentropy");
   1004         let mut session = TangleWebSocketSession::new(
   1005             session_limits(16),
   1006             shutdown.subscribe(),
   1007             runtime,
   1008             auth,
   1009             events,
   1010         )
   1011         .expect("session");
   1012 
   1013         assert_eq!(
   1014             session
   1015                 .dispatch_text("[\"NEG-OPEN\",\"neg-sub\",{\"kinds\":[1]},\"00\"]")
   1016                 .await,
   1017             TangleSessionControl::Continue
   1018         );
   1019         assert_eq!(
   1020             take_outbound_text(&mut session),
   1021             "[\"NEG-ERR\",\"neg-sub\",\"blocked: Negentropy sync is disabled\"]"
   1022         );
   1023         assert_eq!(
   1024             session
   1025                 .dispatch_text("[\"NEG-MSG\",\"neg-sub\",\"\"]")
   1026                 .await,
   1027             TangleSessionControl::Continue
   1028         );
   1029         assert_eq!(
   1030             take_outbound_text(&mut session),
   1031             "[\"NEG-ERR\",\"neg-sub\",\"blocked: Negentropy sync is disabled\"]"
   1032         );
   1033         assert_eq!(
   1034             session.dispatch_text("[\"NEG-CLOSE\",\"neg-sub\"]").await,
   1035             TangleSessionControl::Continue
   1036         );
   1037         assert!(session.outbound_receiver.try_recv().is_err());
   1038     }
   1039 
   1040     #[tokio::test]
   1041     async fn websocket_session_disabled_negentropy_privacy_response_omits_filter_material() {
   1042         let shutdown = TangleShutdownSignal::new();
   1043         let (runtime, auth, events) = session_runtime("disabled-negentropy-privacy");
   1044         let mut session = TangleWebSocketSession::new(
   1045             session_limits(16),
   1046             shutdown.subscribe(),
   1047             runtime,
   1048             auth,
   1049             events,
   1050         )
   1051         .expect("session");
   1052         let hidden_event_id = "a".repeat(64);
   1053         let private_group_id = "private-group-alpha";
   1054         let raw = json!([
   1055             "NEG-OPEN",
   1056             "neg-private",
   1057             {"ids": [hidden_event_id], "#h": [private_group_id]},
   1058             "00"
   1059         ])
   1060         .to_string();
   1061 
   1062         assert_eq!(
   1063             session.dispatch_text(&raw).await,
   1064             TangleSessionControl::Continue
   1065         );
   1066         let text = take_outbound_text(&mut session);
   1067 
   1068         assert_eq!(
   1069             text,
   1070             "[\"NEG-ERR\",\"neg-private\",\"blocked: Negentropy sync is disabled\"]"
   1071         );
   1072         assert!(!text.contains(private_group_id));
   1073         assert!(!text.contains(&hidden_event_id));
   1074         assert!(!text.contains("inventory"));
   1075         assert!(!text.contains("#h"));
   1076     }
   1077 
   1078     #[tokio::test]
   1079     async fn websocket_session_scopes_subscriptions_per_connection() {
   1080         let shutdown = TangleShutdownSignal::new();
   1081         let root = temp_root("connection-scope");
   1082         let _ = std::fs::remove_dir_all(&root);
   1083         let runtime =
   1084             RelayRuntimeHandle::new(RelayRuntime::open(runtime_config(&root)).expect("runtime"));
   1085         let metrics = runtime.metrics();
   1086         let auth_a = runtime.auth_state().await.expect("auth a");
   1087         let auth_b = runtime.auth_state().await.expect("auth b");
   1088         let events_a = runtime.subscribe_events().await;
   1089         let events_b = runtime.subscribe_events().await;
   1090         let mut first = TangleWebSocketSession::new(
   1091             session_limits(8),
   1092             shutdown.subscribe(),
   1093             runtime.clone(),
   1094             auth_a,
   1095             events_a,
   1096         )
   1097         .expect("first");
   1098         let mut second = TangleWebSocketSession::new(
   1099             session_limits(8),
   1100             shutdown.subscribe(),
   1101             runtime.clone(),
   1102             auth_b,
   1103             events_b,
   1104         )
   1105         .expect("second");
   1106         let subscription_id = SubscriptionId::new("shared").expect("subscription");
   1107 
   1108         assert_eq!(
   1109             first
   1110                 .handle_protocol_client_message_for_test(req(subscription_id.clone()))
   1111                 .await
   1112                 .expect("first req"),
   1113             vec![RelayMessage::Eose(subscription_id.clone())]
   1114         );
   1115         assert_eq!(
   1116             second
   1117                 .handle_protocol_client_message_for_test(req(subscription_id.clone()))
   1118                 .await
   1119                 .expect("second req"),
   1120             vec![RelayMessage::Eose(subscription_id.clone())]
   1121         );
   1122         assert_eq!(first.active_subscription_count(), 1);
   1123         assert_eq!(second.active_subscription_count(), 1);
   1124 
   1125         assert_eq!(
   1126             first
   1127                 .handle_protocol_client_message_for_test(ClientMessage::Close(
   1128                     subscription_id.clone()
   1129                 ))
   1130                 .await
   1131                 .expect("close first"),
   1132             Vec::<RelayMessage>::new()
   1133         );
   1134         assert_eq!(first.active_subscription_count(), 0);
   1135         assert_eq!(second.active_subscription_count(), 1);
   1136 
   1137         assert_eq!(
   1138             second
   1139                 .handle_protocol_client_message_for_test(req(subscription_id.clone()))
   1140                 .await
   1141                 .expect("replace second"),
   1142             vec![RelayMessage::Eose(subscription_id.clone())]
   1143         );
   1144         assert_eq!(first.active_subscription_count(), 0);
   1145         assert_eq!(second.active_subscription_count(), 1);
   1146 
   1147         assert_eq!(
   1148             second
   1149                 .handle_protocol_client_message_for_test(ClientMessage::Close(subscription_id))
   1150                 .await
   1151                 .expect("close second"),
   1152             Vec::<RelayMessage>::new()
   1153         );
   1154         assert_eq!(second.active_subscription_count(), 0);
   1155         let snapshot = metrics.snapshot();
   1156         assert_eq!(snapshot.client_messages(), 5);
   1157         assert_eq!(snapshot.req_messages(), 3);
   1158         assert_eq!(snapshot.close_messages(), 2);
   1159         assert_eq!(snapshot.active_subscriptions(), 0);
   1160         assert_eq!(snapshot.opened_subscriptions(), 2);
   1161         assert_eq!(snapshot.closed_subscriptions(), 2);
   1162 
   1163         let _ = std::fs::remove_dir_all(root);
   1164     }
   1165 
   1166     #[tokio::test]
   1167     async fn websocket_session_live_fanout_uses_current_auth() {
   1168         let shutdown = TangleShutdownSignal::new();
   1169         let root = temp_root("current-auth-live");
   1170         let _ = std::fs::remove_dir_all(&root);
   1171         let runtime = RelayRuntimeHandle::new(
   1172             RelayRuntime::open(runtime_config_with_groups(&root)).expect("runtime"),
   1173         );
   1174         let mut owner_auth = runtime.auth_state().await.expect("owner auth");
   1175         owner_auth
   1176             .issue_challenge("owner-live", UnixTimestamp::new(100))
   1177             .expect("owner challenge");
   1178         let owner_auth_event =
   1179             tangle_v2_auth_event(FixtureKey::Owner, "owner-live", 120).expect("owner auth event");
   1180         assert_eq!(
   1181             runtime
   1182                 .handle_protocol_client_message_for_test(
   1183                     ClientMessage::Auth(owner_auth_event.clone()),
   1184                     &mut owner_auth,
   1185                     UnixTimestamp::new(120)
   1186                 )
   1187                 .await
   1188                 .expect("owner auth"),
   1189             vec![RelayMessage::Ok {
   1190                 event_id: owner_auth_event.id().clone(),
   1191                 accepted: true,
   1192                 message: String::new()
   1193             }]
   1194         );
   1195         let create = tangle_v2_group_create_event(FixtureKey::Owner, "LiveFarm", 121, &["private"])
   1196             .expect("create");
   1197         assert_eq!(
   1198             runtime
   1199                 .handle_protocol_client_message_for_test(
   1200                     ClientMessage::Event(create.clone()),
   1201                     &mut owner_auth,
   1202                     UnixTimestamp::new(121)
   1203                 )
   1204                 .await
   1205                 .expect("create"),
   1206             vec![RelayMessage::Ok {
   1207                 event_id: create.id().clone(),
   1208                 accepted: true,
   1209                 message: String::new()
   1210             }]
   1211         );
   1212         let session_auth = runtime.auth_state().await.expect("session auth");
   1213         let events = runtime.subscribe_events().await;
   1214         let mut session = TangleWebSocketSession::new(
   1215             session_limits(8),
   1216             shutdown.subscribe(),
   1217             runtime.clone(),
   1218             session_auth,
   1219             events,
   1220         )
   1221         .expect("session");
   1222         let subscription_id = SubscriptionId::new("current-auth-live").expect("subscription");
   1223 
   1224         assert_eq!(
   1225             session
   1226                 .handle_protocol_client_message_for_test(ClientMessage::Req {
   1227                     subscription_id: subscription_id.clone(),
   1228                     filters: vec![
   1229                         filter_from_value(&json!({"kinds":[1], "#h":["LiveFarm"]}))
   1230                             .expect("filter")
   1231                     ],
   1232                 })
   1233                 .await
   1234                 .expect("req"),
   1235             vec![RelayMessage::Eose(subscription_id.clone())]
   1236         );
   1237         assert_eq!(session.active_subscription_count(), 1);
   1238         let before_auth =
   1239             tangle_v2_group_event(FixtureKey::Owner, "LiveFarm", 122, 1, "before auth")
   1240                 .expect("before auth");
   1241         let before_auth_id = before_auth.id().clone();
   1242         assert_eq!(
   1243             runtime
   1244                 .handle_protocol_client_message_for_test(
   1245                     ClientMessage::Event(before_auth),
   1246                     &mut owner_auth,
   1247                     UnixTimestamp::new(122)
   1248                 )
   1249                 .await
   1250                 .expect("before event"),
   1251             vec![RelayMessage::Ok {
   1252                 event_id: before_auth_id,
   1253                 accepted: true,
   1254                 message: String::new()
   1255             }]
   1256         );
   1257         let offset = session.events.recv().await;
   1258         assert_eq!(
   1259             session.handle_event_receive_result(offset).await,
   1260             TangleSessionControl::Continue
   1261         );
   1262         assert!(session.outbound_receiver.try_recv().is_err());
   1263 
   1264         let session_now = current_unix_timestamp();
   1265         session
   1266             .auth
   1267             .issue_challenge("session-live", session_now)
   1268             .expect("session challenge");
   1269         let session_auth_event =
   1270             tangle_v2_auth_event(FixtureKey::Owner, "session-live", session_now.as_u64())
   1271                 .expect("auth event");
   1272         assert_eq!(
   1273             session
   1274                 .handle_protocol_client_message_for_test(ClientMessage::Auth(
   1275                     session_auth_event.clone()
   1276                 ))
   1277                 .await
   1278                 .expect("session auth"),
   1279             vec![RelayMessage::Ok {
   1280                 event_id: session_auth_event.id().clone(),
   1281                 accepted: true,
   1282                 message: String::new()
   1283             }]
   1284         );
   1285         let after_auth = tangle_v2_group_event(FixtureKey::Owner, "LiveFarm", 132, 1, "after auth")
   1286             .expect("after auth");
   1287         assert_eq!(
   1288             runtime
   1289                 .handle_protocol_client_message_for_test(
   1290                     ClientMessage::Event(after_auth.clone()),
   1291                     &mut owner_auth,
   1292                     UnixTimestamp::new(132)
   1293                 )
   1294                 .await
   1295                 .expect("after event"),
   1296             vec![RelayMessage::Ok {
   1297                 event_id: after_auth.id().clone(),
   1298                 accepted: true,
   1299                 message: String::new()
   1300             }]
   1301         );
   1302         let offset = session.events.recv().await;
   1303         assert_eq!(
   1304             session.handle_event_receive_result(offset).await,
   1305             TangleSessionControl::Continue
   1306         );
   1307         assert_relay_message_text(
   1308             &take_outbound_text(&mut session),
   1309             RelayMessage::Event {
   1310                 subscription_id,
   1311                 event: after_auth,
   1312             },
   1313         );
   1314 
   1315         let _ = std::fs::remove_dir_all(root);
   1316     }
   1317 
   1318     #[tokio::test]
   1319     async fn websocket_session_complete_and_failed_reqs_do_not_subscribe() {
   1320         let shutdown = TangleShutdownSignal::new();
   1321         let root = temp_root("complete-req-lifecycle");
   1322         let _ = std::fs::remove_dir_all(&root);
   1323         let runtime =
   1324             RelayRuntimeHandle::new(RelayRuntime::open(runtime_config(&root)).expect("runtime"));
   1325         let mut auth = runtime.auth_state().await.expect("auth");
   1326         let events = runtime.subscribe_events().await;
   1327         let mut session = TangleWebSocketSession::new(
   1328             session_limits(8),
   1329             shutdown.subscribe(),
   1330             runtime.clone(),
   1331             runtime.auth_state().await.expect("session auth"),
   1332             events,
   1333         )
   1334         .expect("session");
   1335         let event = tangle_v2_event(FixtureKey::Member, 1_714_124_433, 1, Vec::new(), "complete")
   1336             .expect("event");
   1337 
   1338         assert_eq!(
   1339             runtime
   1340                 .handle_protocol_client_message_for_test(
   1341                     ClientMessage::Event(event.clone()),
   1342                     &mut auth,
   1343                     UnixTimestamp::new(1_714_124_433)
   1344                 )
   1345                 .await
   1346                 .expect("event"),
   1347             vec![RelayMessage::Ok {
   1348                 event_id: event.id().clone(),
   1349                 accepted: true,
   1350                 message: String::new()
   1351             }]
   1352         );
   1353         let exact_id = SubscriptionId::new("exact-id").expect("subscription");
   1354         assert_eq!(
   1355             session
   1356                 .handle_protocol_client_message_for_test(ClientMessage::Req {
   1357                     subscription_id: exact_id.clone(),
   1358                     filters: vec![
   1359                         filter_from_value(&json!({"ids":[event.id().as_str()]}))
   1360                             .expect("exact filter")
   1361                     ],
   1362                 })
   1363                 .await
   1364                 .expect("exact req"),
   1365             vec![
   1366                 RelayMessage::Event {
   1367                     subscription_id: exact_id.clone(),
   1368                     event: event.clone()
   1369                 },
   1370                 RelayMessage::Eose(exact_id)
   1371             ]
   1372         );
   1373         assert_eq!(session.active_subscription_count(), 0);
   1374 
   1375         let open = SubscriptionId::new("open").expect("subscription");
   1376         assert_eq!(
   1377             session
   1378                 .handle_protocol_client_message_for_test(ClientMessage::Req {
   1379                     subscription_id: open.clone(),
   1380                     filters: vec![filter_from_value(&json!({"kinds":[1]})).expect("open filter")],
   1381                 })
   1382                 .await
   1383                 .expect("open req"),
   1384             vec![
   1385                 RelayMessage::Event {
   1386                     subscription_id: open.clone(),
   1387                     event
   1388                 },
   1389                 RelayMessage::Eose(open.clone())
   1390             ]
   1391         );
   1392         assert_eq!(session.active_subscription_count(), 1);
   1393 
   1394         let search = SubscriptionId::new("search").expect("subscription");
   1395         assert_eq!(
   1396             session
   1397                 .handle_protocol_client_message_for_test(ClientMessage::Req {
   1398                     subscription_id: search.clone(),
   1399                     filters: vec![
   1400                         filter_from_value(&json!({"search":"carrots"})).expect("search filter")
   1401                     ],
   1402                 })
   1403                 .await
   1404                 .expect("search req"),
   1405             vec![RelayMessage::Closed {
   1406                 subscription_id: search,
   1407                 message: "unsupported: search filters are not supported".to_owned()
   1408             }]
   1409         );
   1410         assert_eq!(session.active_subscription_count(), 1);
   1411 
   1412         let invalid = SubscriptionId::new("invalid").expect("subscription");
   1413         let invalid_result = session
   1414             .handle_protocol_client_message_for_test(ClientMessage::Req {
   1415                 subscription_id: invalid,
   1416                 filters: vec![Filter::empty(); 11],
   1417             })
   1418             .await;
   1419         assert!(invalid_result.is_err());
   1420         assert_eq!(session.active_subscription_count(), 1);
   1421 
   1422         let _ = std::fs::remove_dir_all(root);
   1423     }
   1424 
   1425     #[tokio::test]
   1426     async fn websocket_session_redacted_initial_req_closes_without_live_subscription() {
   1427         let shutdown = TangleShutdownSignal::new();
   1428         let root = temp_root("redacted-req-close");
   1429         let _ = std::fs::remove_dir_all(&root);
   1430         let runtime = RelayRuntimeHandle::new(
   1431             RelayRuntime::open(runtime_config_with_groups(&root)).expect("runtime"),
   1432         );
   1433         let mut owner_auth = runtime.auth_state().await.expect("owner auth");
   1434         owner_auth
   1435             .issue_challenge("owner-redacted", UnixTimestamp::new(100))
   1436             .expect("owner challenge");
   1437         let owner_auth_event = tangle_v2_auth_event(FixtureKey::Owner, "owner-redacted", 120)
   1438             .expect("owner auth event");
   1439         assert_eq!(
   1440             runtime
   1441                 .handle_protocol_client_message_for_test(
   1442                     ClientMessage::Auth(owner_auth_event.clone()),
   1443                     &mut owner_auth,
   1444                     UnixTimestamp::new(120)
   1445                 )
   1446                 .await
   1447                 .expect("owner auth"),
   1448             vec![RelayMessage::Ok {
   1449                 event_id: owner_auth_event.id().clone(),
   1450                 accepted: true,
   1451                 message: String::new()
   1452             }]
   1453         );
   1454         let create =
   1455             tangle_v2_group_create_event(FixtureKey::Owner, "RedactedFarm", 121, &["private"])
   1456                 .expect("create");
   1457         assert_eq!(
   1458             runtime
   1459                 .handle_protocol_client_message_for_test(
   1460                     ClientMessage::Event(create.clone()),
   1461                     &mut owner_auth,
   1462                     UnixTimestamp::new(121)
   1463                 )
   1464                 .await
   1465                 .expect("create"),
   1466             vec![RelayMessage::Ok {
   1467                 event_id: create.id().clone(),
   1468                 accepted: true,
   1469                 message: String::new()
   1470             }]
   1471         );
   1472         let public_event =
   1473             tangle_v2_event(FixtureKey::Member, 1_714_124_433, 1, Vec::new(), "public")
   1474                 .expect("public");
   1475         assert_eq!(
   1476             runtime
   1477                 .handle_protocol_client_message_for_test(
   1478                     ClientMessage::Event(public_event.clone()),
   1479                     &mut owner_auth,
   1480                     UnixTimestamp::new(122)
   1481                 )
   1482                 .await
   1483                 .expect("public event"),
   1484             vec![RelayMessage::Ok {
   1485                 event_id: public_event.id().clone(),
   1486                 accepted: true,
   1487                 message: String::new()
   1488             }]
   1489         );
   1490         let private_event =
   1491             tangle_v2_group_event(FixtureKey::Owner, "RedactedFarm", 123, 1, "private")
   1492                 .expect("private");
   1493         assert_eq!(
   1494             runtime
   1495                 .handle_protocol_client_message_for_test(
   1496                     ClientMessage::Event(private_event.clone()),
   1497                     &mut owner_auth,
   1498                     UnixTimestamp::new(123)
   1499                 )
   1500                 .await
   1501                 .expect("private event"),
   1502             vec![RelayMessage::Ok {
   1503                 event_id: private_event.id().clone(),
   1504                 accepted: true,
   1505                 message: String::new()
   1506             }]
   1507         );
   1508 
   1509         let events = runtime.subscribe_events().await;
   1510         let mut session = TangleWebSocketSession::new(
   1511             session_limits(8),
   1512             shutdown.subscribe(),
   1513             runtime.clone(),
   1514             runtime.auth_state().await.expect("session auth"),
   1515             events,
   1516         )
   1517         .expect("session");
   1518         let subscription_id = SubscriptionId::new("redacted-req").expect("subscription");
   1519         assert_eq!(
   1520             session
   1521                 .handle_protocol_client_message_for_test(ClientMessage::Req {
   1522                     subscription_id: subscription_id.clone(),
   1523                     filters: vec![filter_from_value(&json!({"kinds":[1]})).expect("filter")],
   1524                 })
   1525                 .await
   1526                 .expect("redacted req"),
   1527             vec![
   1528                 RelayMessage::Event {
   1529                     subscription_id: subscription_id.clone(),
   1530                     event: public_event
   1531                 },
   1532                 RelayMessage::Closed {
   1533                     subscription_id,
   1534                     message: "auth-required: authentication required to read group events"
   1535                         .to_owned()
   1536                 }
   1537             ]
   1538         );
   1539         assert_eq!(session.active_subscription_count(), 0);
   1540 
   1541         let _ = std::fs::remove_dir_all(root);
   1542     }
   1543 
   1544     #[tokio::test]
   1545     async fn websocket_session_preserves_chorus_close_scope_parity() {
   1546         let shutdown = TangleShutdownSignal::new();
   1547         let root = temp_root("chorus-close-scope-parity");
   1548         let _ = std::fs::remove_dir_all(&root);
   1549         let runtime =
   1550             RelayRuntimeHandle::new(RelayRuntime::open(runtime_config(&root)).expect("runtime"));
   1551         let metrics = runtime.metrics();
   1552         let auth_a = runtime.auth_state().await.expect("auth a");
   1553         let auth_b = runtime.auth_state().await.expect("auth b");
   1554         let events_a = runtime.subscribe_events().await;
   1555         let events_b = runtime.subscribe_events().await;
   1556         let mut first = TangleWebSocketSession::new(
   1557             session_limits(8),
   1558             shutdown.subscribe(),
   1559             runtime.clone(),
   1560             auth_a,
   1561             events_a,
   1562         )
   1563         .expect("first");
   1564         let mut second = TangleWebSocketSession::new(
   1565             session_limits(8),
   1566             shutdown.subscribe(),
   1567             runtime,
   1568             auth_b,
   1569             events_b,
   1570         )
   1571         .expect("second");
   1572         let subscription_id = SubscriptionId::new("shared-close").expect("subscription");
   1573         let req_text = json!(["REQ", subscription_id.as_str(), {"kinds":[1]}]).to_string();
   1574 
   1575         assert_eq!(
   1576             first.dispatch_text(&req_text).await,
   1577             TangleSessionControl::Continue
   1578         );
   1579         assert_eq!(
   1580             take_outbound_text(&mut first),
   1581             RelayMessage::Eose(subscription_id.clone()).encode()
   1582         );
   1583         assert_eq!(
   1584             second.dispatch_text(&req_text).await,
   1585             TangleSessionControl::Continue
   1586         );
   1587         assert_eq!(
   1588             take_outbound_text(&mut second),
   1589             RelayMessage::Eose(subscription_id.clone()).encode()
   1590         );
   1591         assert_eq!(first.active_subscription_count(), 1);
   1592         assert_eq!(second.active_subscription_count(), 1);
   1593 
   1594         let close_text = json!(["CLOSE", subscription_id.as_str()]).to_string();
   1595         assert_eq!(
   1596             first.dispatch_text(&close_text).await,
   1597             TangleSessionControl::Continue
   1598         );
   1599         assert!(first.outbound_receiver.try_recv().is_err());
   1600         assert_eq!(
   1601             first.dispatch_text(&close_text).await,
   1602             TangleSessionControl::Continue
   1603         );
   1604         assert!(first.outbound_receiver.try_recv().is_err());
   1605         assert_eq!(first.active_subscription_count(), 0);
   1606         assert_eq!(second.active_subscription_count(), 1);
   1607 
   1608         let event = tangle_v2_event(
   1609             FixtureKey::Member,
   1610             1_714_124_433,
   1611             1,
   1612             Vec::new(),
   1613             "close scope parity",
   1614         )
   1615         .expect("event");
   1616         assert_eq!(
   1617             first
   1618                 .dispatch_text(&json!(["EVENT", event_to_value(&event)]).to_string())
   1619                 .await,
   1620             TangleSessionControl::Continue
   1621         );
   1622         assert_eq!(
   1623             take_outbound_text(&mut first),
   1624             RelayMessage::Ok {
   1625                 event_id: event.id().clone(),
   1626                 accepted: true,
   1627                 message: String::new()
   1628             }
   1629             .encode()
   1630         );
   1631 
   1632         let first_offset = first.events.recv().await;
   1633         let second_offset = second.events.recv().await;
   1634         assert_eq!(
   1635             first.handle_event_receive_result(first_offset).await,
   1636             TangleSessionControl::Continue
   1637         );
   1638         assert!(first.outbound_receiver.try_recv().is_err());
   1639         assert_eq!(
   1640             second.handle_event_receive_result(second_offset).await,
   1641             TangleSessionControl::Continue
   1642         );
   1643         assert_relay_message_text(
   1644             &take_outbound_text(&mut second),
   1645             RelayMessage::Event {
   1646                 subscription_id: subscription_id.clone(),
   1647                 event,
   1648             },
   1649         );
   1650         let snapshot = metrics.snapshot();
   1651         assert_eq!(snapshot.client_messages(), 5);
   1652         assert_eq!(snapshot.event_messages(), 1);
   1653         assert_eq!(snapshot.req_messages(), 2);
   1654         assert_eq!(snapshot.close_messages(), 2);
   1655         assert_eq!(snapshot.opened_subscriptions(), 2);
   1656         assert_eq!(snapshot.closed_subscriptions(), 1);
   1657 
   1658         let _ = std::fs::remove_dir_all(root);
   1659     }
   1660 
   1661     #[tokio::test]
   1662     async fn websocket_session_rate_limited_req_does_not_subscribe() {
   1663         let shutdown = TangleShutdownSignal::new();
   1664         let root = temp_root("rate-limited-req");
   1665         let _ = std::fs::remove_dir_all(&root);
   1666         let runtime = RelayRuntime::open(runtime_config(&root)).expect("runtime");
   1667         let rule = runtime.config().rate_limits().req().per_connection();
   1668         let runtime = RelayRuntimeHandle::new(runtime);
   1669         let auth = runtime.auth_state().await.expect("auth");
   1670         let events = runtime.subscribe_events().await;
   1671         let now = current_unix_timestamp();
   1672         let mut session = TangleWebSocketSession::new(
   1673             session_limits(8),
   1674             shutdown.subscribe(),
   1675             runtime.clone(),
   1676             auth,
   1677             events,
   1678         )
   1679         .expect("session");
   1680         let key = TangleRateLimitKey::connection(TangleRateLimitScope::Req, session.connection_id);
   1681         let limiter = runtime.rate_limiter().await;
   1682         for _ in 0..rule.max_hits() {
   1683             limiter.record(key.clone(), rule, now);
   1684         }
   1685         let subscription_id = SubscriptionId::new("limited").expect("subscription");
   1686 
   1687         assert_eq!(
   1688             session
   1689                 .handle_protocol_client_message_for_test(ClientMessage::Req {
   1690                     subscription_id: subscription_id.clone(),
   1691                     filters: vec![
   1692                         filter_from_value(&json!({"kinds": [1], "limit": 1})).expect("filter")
   1693                     ]
   1694                 })
   1695                 .await
   1696                 .expect("req"),
   1697             vec![RelayMessage::Closed {
   1698                 subscription_id,
   1699                 message: format!(
   1700                     "rate-limited: req connection rate limit exceeded until {}",
   1701                     now.as_u64() + 60
   1702                 )
   1703             }]
   1704         );
   1705         assert_eq!(session.active_subscription_count(), 0);
   1706         let snapshot = runtime.metrics().snapshot();
   1707         assert_eq!(snapshot.client_messages(), 1);
   1708         assert_eq!(snapshot.req_messages(), 1);
   1709         assert_eq!(snapshot.opened_subscriptions(), 0);
   1710         assert_eq!(snapshot.rate_limit_rejections(), 1);
   1711 
   1712         let _ = std::fs::remove_dir_all(root);
   1713     }
   1714 
   1715     #[tokio::test]
   1716     async fn websocket_session_closes_when_event_receiver_lags() {
   1717         let shutdown = TangleShutdownSignal::new();
   1718         let root = temp_root("event-receiver-lag");
   1719         let _ = std::fs::remove_dir_all(&root);
   1720         let runtime =
   1721             RelayRuntime::open(runtime_config_with_outbound_queue(&root, 1)).expect("runtime");
   1722         let auth = runtime.auth_state().expect("auth");
   1723         let events = runtime.event_bus().subscribe();
   1724         assert_eq!(runtime.event_bus().publish(StoreOffset::new(1)), 1);
   1725         assert_eq!(runtime.event_bus().publish(StoreOffset::new(2)), 1);
   1726         let runtime = RelayRuntimeHandle::new(runtime);
   1727         let metrics = runtime.metrics();
   1728         let mut session = TangleWebSocketSession::new(
   1729             session_limits(1),
   1730             shutdown.subscribe(),
   1731             runtime,
   1732             auth,
   1733             events,
   1734         )
   1735         .expect("session");
   1736         let event = session.events.recv().await;
   1737 
   1738         assert_eq!(
   1739             session.handle_event_receive_result(event).await,
   1740             TangleSessionControl::Close(event_stream_lag_close_message())
   1741         );
   1742         assert_eq!(metrics.event_bus_lagged_receivers(), 1);
   1743         assert_eq!(metrics.event_bus_lagged_offsets(), 1);
   1744 
   1745         let _ = std::fs::remove_dir_all(root);
   1746     }
   1747 
   1748     #[tokio::test]
   1749     async fn websocket_session_preserves_chorus_live_fanout_backpressure_parity() {
   1750         let shutdown = TangleShutdownSignal::new();
   1751         let live_root = temp_root("chorus-live-fanout-parity");
   1752         let _ = std::fs::remove_dir_all(&live_root);
   1753         let runtime = RelayRuntimeHandle::new(
   1754             RelayRuntime::open(runtime_config_with_outbound_queue(&live_root, 1)).expect("runtime"),
   1755         );
   1756         let metrics = runtime.metrics();
   1757         let auth = runtime.auth_state().await.expect("auth");
   1758         let events = runtime.subscribe_events().await;
   1759         let mut session = TangleWebSocketSession::new(
   1760             session_limits(1),
   1761             shutdown.subscribe(),
   1762             runtime,
   1763             auth,
   1764             events,
   1765         )
   1766         .expect("session");
   1767         let subscription_id = SubscriptionId::new("chorus-live").expect("subscription");
   1768         let req_text = json!(["REQ", subscription_id.as_str(), {"kinds":[1]}]).to_string();
   1769 
   1770         assert_eq!(
   1771             session.dispatch_text(&req_text).await,
   1772             TangleSessionControl::Continue
   1773         );
   1774         assert_eq!(
   1775             take_outbound_text(&mut session),
   1776             RelayMessage::Eose(subscription_id.clone()).encode()
   1777         );
   1778         for index in 0..3 {
   1779             let content = format!("chorus live {index}");
   1780             let event = tangle_v2_event(
   1781                 FixtureKey::Member,
   1782                 1_714_124_433 + index,
   1783                 1,
   1784                 Vec::new(),
   1785                 &content,
   1786             )
   1787             .expect("event");
   1788             assert_eq!(
   1789                 session
   1790                     .dispatch_text(&json!(["EVENT", event_to_value(&event)]).to_string())
   1791                     .await,
   1792                 TangleSessionControl::Continue
   1793             );
   1794             assert_eq!(
   1795                 take_outbound_text(&mut session),
   1796                 RelayMessage::Ok {
   1797                     event_id: event.id().clone(),
   1798                     accepted: true,
   1799                     message: String::new()
   1800                 }
   1801                 .encode()
   1802             );
   1803             let offset = session.events.recv().await;
   1804             assert_eq!(
   1805                 session.handle_event_receive_result(offset).await,
   1806                 TangleSessionControl::Continue
   1807             );
   1808             assert_relay_message_text(
   1809                 &take_outbound_text(&mut session),
   1810                 RelayMessage::Event {
   1811                     subscription_id: subscription_id.clone(),
   1812                     event,
   1813                 },
   1814             );
   1815             assert_eq!(session.active_subscription_count(), 1);
   1816         }
   1817         assert_eq!(metrics.outbound_queue_full_closes(), 0);
   1818         assert_eq!(metrics.event_bus_lagged_receivers(), 0);
   1819         assert_eq!(metrics.event_bus_lagged_offsets(), 0);
   1820         let _ = std::fs::remove_dir_all(live_root);
   1821 
   1822         let lag_root = temp_root("chorus-live-lag-parity");
   1823         let _ = std::fs::remove_dir_all(&lag_root);
   1824         let runtime =
   1825             RelayRuntime::open(runtime_config_with_outbound_queue(&lag_root, 1)).expect("runtime");
   1826         let auth = runtime.auth_state().expect("auth");
   1827         let events = runtime.event_bus().subscribe();
   1828         assert_eq!(runtime.event_bus().publish(StoreOffset::new(1)), 1);
   1829         assert_eq!(runtime.event_bus().publish(StoreOffset::new(2)), 1);
   1830         let runtime = RelayRuntimeHandle::new(runtime);
   1831         let metrics = runtime.metrics();
   1832         let mut lagged = TangleWebSocketSession::new(
   1833             session_limits(1),
   1834             shutdown.subscribe(),
   1835             runtime,
   1836             auth,
   1837             events,
   1838         )
   1839         .expect("lagged");
   1840         let event = lagged.events.recv().await;
   1841         assert_eq!(
   1842             lagged.handle_event_receive_result(event).await,
   1843             TangleSessionControl::Close(event_stream_lag_close_message())
   1844         );
   1845         assert_eq!(metrics.event_bus_lagged_receivers(), 1);
   1846         assert_eq!(metrics.event_bus_lagged_offsets(), 1);
   1847         let _ = std::fs::remove_dir_all(lag_root);
   1848 
   1849         let (runtime, auth, events) = session_runtime("chorus-outbound-full-parity");
   1850         let metrics = runtime.metrics();
   1851         let mut blocked = TangleWebSocketSession::new(
   1852             session_limits(1),
   1853             shutdown.subscribe(),
   1854             runtime,
   1855             auth,
   1856             events,
   1857         )
   1858         .expect("blocked");
   1859         blocked
   1860             .outbound()
   1861             .try_send(Message::Text("blocked".into()))
   1862             .expect("fill queue");
   1863         assert_eq!(
   1864             blocked.dispatch_text("{").await,
   1865             TangleSessionControl::Close(outbound_queue_full_close_message())
   1866         );
   1867         assert_eq!(metrics.outbound_queue_full_closes(), 1);
   1868     }
   1869 
   1870     #[test]
   1871     fn outbound_queue_is_bounded() {
   1872         let shutdown = TangleShutdownSignal::new();
   1873         let (runtime, auth, events) = session_runtime("outbound-queue");
   1874         let session = TangleWebSocketSession::new(
   1875             session_limits(1),
   1876             shutdown.subscribe(),
   1877             runtime,
   1878             auth,
   1879             events,
   1880         )
   1881         .expect("session");
   1882         let outbound = session.outbound();
   1883 
   1884         assert_eq!(outbound.capacity(), 1);
   1885         outbound
   1886             .try_send(Message::Text("first".into()))
   1887             .expect("first");
   1888         assert_eq!(
   1889             outbound
   1890                 .try_send(Message::Text("second".into()))
   1891                 .expect_err("full"),
   1892             TangleOutboundQueueError::Full
   1893         );
   1894     }
   1895 
   1896     #[tokio::test]
   1897     async fn websocket_session_closes_when_outbound_queue_is_full() {
   1898         let shutdown = TangleShutdownSignal::new();
   1899         let (runtime, auth, events) = session_runtime("outbound-queue-full-close");
   1900         let metrics = runtime.metrics();
   1901         let mut session = TangleWebSocketSession::new(
   1902             session_limits(1),
   1903             shutdown.subscribe(),
   1904             runtime,
   1905             auth,
   1906             events,
   1907         )
   1908         .expect("session");
   1909         session
   1910             .outbound()
   1911             .try_send(Message::Text("blocked".into()))
   1912             .expect("fill queue");
   1913 
   1914         assert_eq!(
   1915             session.dispatch_text("{").await,
   1916             TangleSessionControl::Close(outbound_queue_full_close_message())
   1917         );
   1918         assert_eq!(metrics.outbound_queue_full_closes(), 1);
   1919     }
   1920 
   1921     fn tangle_v2_event(
   1922         key: FixtureKey,
   1923         created_at: u64,
   1924         kind: u64,
   1925         tags: Vec<Tag>,
   1926         content: &str,
   1927     ) -> Result<Event, String> {
   1928         let event = session_pocket_event(key, created_at, kind, tags, content);
   1929         session_pocket_event_to_protocol(&event)
   1930     }
   1931 
   1932     fn tangle_v2_auth_event(
   1933         key: FixtureKey,
   1934         challenge: &str,
   1935         created_at: u64,
   1936     ) -> Result<Event, String> {
   1937         tangle_v2_event(
   1938             key,
   1939             created_at,
   1940             22_242,
   1941             vec![
   1942                 Tag::from_parts("relay", &["wss://relay.radroots.test"])?,
   1943                 Tag::from_parts("challenge", &[challenge])?,
   1944             ],
   1945             "",
   1946         )
   1947     }
   1948 
   1949     fn tangle_v2_group_create_event(
   1950         key: FixtureKey,
   1951         group_id: &str,
   1952         created_at: u64,
   1953         flags: &[&str],
   1954     ) -> Result<Event, String> {
   1955         let mut tags = vec![
   1956             Tag::from_parts("h", &[group_id])?,
   1957             Tag::from_parts("name", &[group_id])?,
   1958         ];
   1959         for flag in flags {
   1960             tags.push(Tag::from_parts(flag, &[])?);
   1961         }
   1962         tangle_v2_event(key, created_at, KIND_GROUP_CREATE_GROUP.into(), tags, "")
   1963     }
   1964 
   1965     fn tangle_v2_group_event(
   1966         key: FixtureKey,
   1967         group_id: &str,
   1968         created_at: u64,
   1969         kind: u64,
   1970         content: &str,
   1971     ) -> Result<Event, String> {
   1972         tangle_v2_event(
   1973             key,
   1974             created_at,
   1975             kind,
   1976             vec![Tag::from_parts("h", &[group_id])?],
   1977             content,
   1978         )
   1979     }
   1980 
   1981     fn session_pocket_event(
   1982         key: FixtureKey,
   1983         created_at: u64,
   1984         kind: u64,
   1985         tags: Vec<Tag>,
   1986         content: &str,
   1987     ) -> PocketOwnedEvent {
   1988         let tags = session_pocket_tags_from_protocol(&tags);
   1989         let secret = format!("{:02x}", fixture_secret_byte(key)).repeat(32);
   1990         RelaySigner::from_secret_hex(&secret)
   1991             .expect("signer")
   1992             .sign_pocket_event(
   1993                 PocketKind::from_u16(u16::try_from(kind).expect("pocket kind")),
   1994                 &tags,
   1995                 PocketTime::from_u64(created_at),
   1996                 content.as_bytes(),
   1997             )
   1998             .expect("pocket event")
   1999     }
   2000 
   2001     fn session_pocket_tags_from_protocol(tags: &[Tag]) -> PocketOwnedTags {
   2002         let parts = tags
   2003             .iter()
   2004             .map(|tag| tag.values().iter().map(String::as_str).collect::<Vec<_>>())
   2005             .collect::<Vec<_>>();
   2006         PocketOwnedTags::new(&parts).expect("pocket tags")
   2007     }
   2008 
   2009     fn session_pocket_event_to_protocol(event: &PocketEvent) -> Result<Event, String> {
   2010         let tags = event
   2011             .tags()
   2012             .map_err(|error| error.to_string())?
   2013             .iter()
   2014             .map(|tag| {
   2015                 Tag::new(
   2016                     tag.map(|value| {
   2017                         std::str::from_utf8(value)
   2018                             .map(str::to_owned)
   2019                             .map_err(|error| error.to_string())
   2020                     })
   2021                     .collect::<Result<Vec<_>, _>>()?,
   2022                 )
   2023                 .map_err(|error| error.to_string())
   2024             })
   2025             .collect::<Result<Vec<_>, _>>()?;
   2026         Ok(Event::new(
   2027             EventId::new(&event.id().as_hex_string()).map_err(|error| error.to_string())?,
   2028             UnsignedEvent::new(
   2029                 PublicKeyHex::new(&event.pubkey().as_hex_string())
   2030                     .map_err(|error| error.to_string())?,
   2031                 UnixTimestamp::new(event.created_at().as_u64()),
   2032                 Kind::new(u64::from(event.kind().as_u16())).map_err(|error| error.to_string())?,
   2033                 tags,
   2034                 std::str::from_utf8(event.content()).map_err(|error| error.to_string())?,
   2035             ),
   2036             SignatureHex::new(&event.sig().to_string()).map_err(|error| error.to_string())?,
   2037         ))
   2038     }
   2039 
   2040     fn fixture_secret_byte(key: FixtureKey) -> u8 {
   2041         match key {
   2042             FixtureKey::Relay => 9,
   2043             FixtureKey::Owner => 10,
   2044             FixtureKey::Admin => 11,
   2045             FixtureKey::Member => 12,
   2046             FixtureKey::Outsider => 13,
   2047         }
   2048     }
   2049 
   2050     fn session_runtime(
   2051         name: &str,
   2052     ) -> (
   2053         RelayRuntimeHandle,
   2054         crate::relay::auth::BaseAuthState,
   2055         TangleEventReceiver,
   2056     ) {
   2057         let root = temp_root(name);
   2058         let _ = std::fs::remove_dir_all(&root);
   2059         let runtime = RelayRuntime::open(runtime_config(&root)).expect("runtime");
   2060         let auth = runtime.auth_state().expect("auth");
   2061         let events = runtime.event_bus().subscribe();
   2062         (RelayRuntimeHandle::new(runtime), auth, events)
   2063     }
   2064 
   2065     fn req(subscription_id: SubscriptionId) -> ClientMessage {
   2066         ClientMessage::Req {
   2067             subscription_id,
   2068             filters: vec![Filter::empty()],
   2069         }
   2070     }
   2071 
   2072     fn take_outbound_text(session: &mut TangleWebSocketSession) -> String {
   2073         let message = session.outbound_receiver.try_recv().expect("message");
   2074         let Message::Text(text) = message else {
   2075             panic!("expected text message")
   2076         };
   2077         text.to_string()
   2078     }
   2079 
   2080     fn assert_relay_message_text(actual: &str, expected: RelayMessage) {
   2081         assert_eq!(
   2082             serde_json::from_str::<serde_json::Value>(actual).expect("actual relay JSON"),
   2083             serde_json::from_str::<serde_json::Value>(&expected.encode())
   2084                 .expect("expected relay JSON")
   2085         );
   2086     }
   2087 
   2088     fn runtime_config(root: &Path) -> BaseRelayRuntimeConfig {
   2089         runtime_config_with_outbound_queue(root, 8)
   2090     }
   2091 
   2092     fn runtime_config_with_groups(root: &Path) -> BaseRelayRuntimeConfig {
   2093         let raw = json!({
   2094             "server": {
   2095                 "listen_addr": "127.0.0.1:0",
   2096                 "relay_url": "wss://relay.radroots.test"
   2097             },
   2098             "pocket": {
   2099                 "data_directory": root.join("pocket"),
   2100                 "sync_policy": "flush_on_shutdown",
   2101                 "query": {
   2102                   "allow_scraping": false,
   2103                   "allow_scrape_if_limited_to": 100,
   2104                   "allow_scrape_if_max_seconds": 3600
   2105                 }
   2106             },
   2107             "groups": {
   2108                 "enabled": true,
   2109                 "canonical_relay_url": "wss://relay.radroots.test",
   2110                 "relay_secret": "7777777777777777777777777777777777777777777777777777777777777777",
   2111                 "owner_pubkeys": [FixtureKey::Owner.public_key().as_str()],
   2112                 "policy": {
   2113                     "public_join": false,
   2114                     "invites_enabled": false
   2115                 }
   2116             },
   2117             "auth": {
   2118                 "challenge_ttl_seconds": 300,
   2119                 "created_at_skew_seconds": 600
   2120             },
   2121             "limits": {
   2122                 "max_message_length": 1048576,
   2123                 "max_subid_length": 64,
   2124                 "max_subscriptions_per_connection": 64,
   2125                 "max_filters_per_request": 10,
   2126                 "max_tag_values_per_filter": 100,
   2127                 "max_query_complexity": 2048,
   2128                 "max_limit": 500,
   2129                 "default_limit": 100,
   2130                 "max_event_tags": 200,
   2131                 "max_content_length": 65536,
   2132                 "broadcast_channel_capacity": 8,
   2133                 "per_connection_outbound_queue": 8
   2134             },
   2135             "rate_limits": {
   2136                 "auth": {
   2137                     "per_ip": {"window_seconds": 60, "max_hits": 120},
   2138                     "per_pubkey": {"window_seconds": 60, "max_hits": 30},
   2139                     "failures": {"window_seconds": 300, "max_hits": 5},
   2140                     "failures_per_ip": {"window_seconds": 300, "max_hits": 20}
   2141                 },
   2142                 "event": {
   2143                     "per_ip": {"window_seconds": 60, "max_hits": 600},
   2144                     "per_pubkey": {"window_seconds": 60, "max_hits": 120},
   2145                     "per_kind": {"window_seconds": 60, "max_hits": 1000}
   2146                 },
   2147                 "group": {
   2148                     "write_per_ip": {"window_seconds": 60, "max_hits": 300},
   2149                     "write_per_pubkey": {"window_seconds": 60, "max_hits": 60},
   2150                     "write_per_group": {"window_seconds": 60, "max_hits": 90},
   2151                     "write_per_kind": {"window_seconds": 60, "max_hits": 300},
   2152                     "join_flow": {"window_seconds": 300, "max_hits": 10},
   2153                     "join_flow_per_ip": {"window_seconds": 300, "max_hits": 30}
   2154                 },
   2155                 "req": {
   2156                     "per_ip": {"window_seconds": 60, "max_hits": 600},
   2157                     "per_connection": {"window_seconds": 60, "max_hits": 120},
   2158                     "per_pubkey": {"window_seconds": 60, "max_hits": 240},
   2159                     "per_group": {"window_seconds": 60, "max_hits": 240},
   2160                     "per_kind": {"window_seconds": 60, "max_hits": 500},
   2161                     "broad": {"window_seconds": 60, "max_hits": 30}
   2162                 },
   2163                 "count": {
   2164                     "per_ip": {"window_seconds": 60, "max_hits": 300},
   2165                     "per_connection": {"window_seconds": 60, "max_hits": 60},
   2166                     "per_pubkey": {"window_seconds": 60, "max_hits": 120},
   2167                     "per_group": {"window_seconds": 60, "max_hits": 120},
   2168                     "per_kind": {"window_seconds": 60, "max_hits": 240},
   2169                     "broad": {"window_seconds": 60, "max_hits": 20}
   2170                 }
   2171             }
   2172         })
   2173         .to_string();
   2174         parse_base_relay_runtime_config_json(&raw).expect("config")
   2175     }
   2176 
   2177     fn runtime_config_with_outbound_queue(
   2178         root: &Path,
   2179         per_connection_outbound_queue: usize,
   2180     ) -> BaseRelayRuntimeConfig {
   2181         let raw = json!({
   2182             "server": {
   2183                 "listen_addr": "127.0.0.1:0",
   2184                 "relay_url": "wss://relay.radroots.test"
   2185             },
   2186             "pocket": {
   2187                 "data_directory": root.join("pocket"),
   2188                 "sync_policy": "flush_on_shutdown",
   2189                 "query": {
   2190                   "allow_scraping": false,
   2191                   "allow_scrape_if_limited_to": 100,
   2192                   "allow_scrape_if_max_seconds": 3600
   2193                 }
   2194             },
   2195             "groups": {
   2196                 "enabled": false
   2197             },
   2198             "auth": {
   2199                 "challenge_ttl_seconds": 300,
   2200                 "created_at_skew_seconds": 600
   2201             },
   2202             "limits": {
   2203                 "max_message_length": 1048576,
   2204                 "max_subid_length": 64,
   2205                 "max_subscriptions_per_connection": 64,
   2206                 "max_filters_per_request": 10,
   2207                 "max_tag_values_per_filter": 100,
   2208                 "max_query_complexity": 2048,
   2209                 "max_limit": 500,
   2210                 "default_limit": 100,
   2211                 "max_event_tags": 200,
   2212                 "max_content_length": 65536,
   2213                 "broadcast_channel_capacity": per_connection_outbound_queue,
   2214                 "per_connection_outbound_queue": per_connection_outbound_queue
   2215             },
   2216             "rate_limits": {
   2217                 "auth": {
   2218                     "per_ip": {"window_seconds": 60, "max_hits": 120},
   2219                     "per_pubkey": {"window_seconds": 60, "max_hits": 30},
   2220                     "failures": {"window_seconds": 300, "max_hits": 5},
   2221                     "failures_per_ip": {"window_seconds": 300, "max_hits": 20}
   2222                 },
   2223                 "event": {
   2224                     "per_ip": {"window_seconds": 60, "max_hits": 600},
   2225                     "per_pubkey": {"window_seconds": 60, "max_hits": 120},
   2226                     "per_kind": {"window_seconds": 60, "max_hits": 1000}
   2227                 },
   2228                 "group": {
   2229                     "write_per_ip": {"window_seconds": 60, "max_hits": 300},
   2230                     "write_per_pubkey": {"window_seconds": 60, "max_hits": 60},
   2231                     "write_per_group": {"window_seconds": 60, "max_hits": 90},
   2232                     "write_per_kind": {"window_seconds": 60, "max_hits": 300},
   2233                     "join_flow": {"window_seconds": 300, "max_hits": 10},
   2234                     "join_flow_per_ip": {"window_seconds": 300, "max_hits": 30}
   2235                 },
   2236                 "req": {
   2237                     "per_ip": {"window_seconds": 60, "max_hits": 600},
   2238                     "per_connection": {"window_seconds": 60, "max_hits": 120},
   2239                     "per_pubkey": {"window_seconds": 60, "max_hits": 240},
   2240                     "per_group": {"window_seconds": 60, "max_hits": 240},
   2241                     "per_kind": {"window_seconds": 60, "max_hits": 500},
   2242                     "broad": {"window_seconds": 60, "max_hits": 30}
   2243                 },
   2244                 "count": {
   2245                     "per_ip": {"window_seconds": 60, "max_hits": 300},
   2246                     "per_connection": {"window_seconds": 60, "max_hits": 60},
   2247                     "per_pubkey": {"window_seconds": 60, "max_hits": 120},
   2248                     "per_group": {"window_seconds": 60, "max_hits": 120},
   2249                     "per_kind": {"window_seconds": 60, "max_hits": 240},
   2250                     "broad": {"window_seconds": 60, "max_hits": 20}
   2251                 }
   2252             }
   2253         })
   2254         .to_string();
   2255         parse_base_relay_runtime_config_json(&raw).expect("config")
   2256     }
   2257 
   2258     fn session_limits(per_connection_outbound_queue: usize) -> TangleRuntimeLimits {
   2259         session_limits_result(per_connection_outbound_queue).expect("limits")
   2260     }
   2261 
   2262     fn session_limits_with_message_length(
   2263         max_message_length: usize,
   2264         per_connection_outbound_queue: usize,
   2265     ) -> TangleRuntimeLimits {
   2266         TangleRuntimeLimits::new(
   2267             max_message_length,
   2268             BaseRelayLimits::new(BaseRelayLimitSettings {
   2269                 max_pending_events: per_connection_outbound_queue,
   2270                 max_subscription_id_length: 64,
   2271                 max_subscriptions: 64,
   2272                 max_filters_per_request: 10,
   2273                 max_tag_values_per_filter: 100,
   2274                 max_query_complexity: 610,
   2275                 max_event_tags: 200,
   2276                 max_content_length: 65_536,
   2277                 max_limit: 500,
   2278                 default_limit: 100,
   2279             })
   2280             .expect("relay limits"),
   2281             16,
   2282             per_connection_outbound_queue,
   2283         )
   2284         .expect("limits")
   2285     }
   2286 
   2287     fn session_limits_result(
   2288         per_connection_outbound_queue: usize,
   2289     ) -> Result<TangleRuntimeLimits, BaseRelayError> {
   2290         TangleRuntimeLimits::new(
   2291             1_048_576,
   2292             BaseRelayLimits::new(BaseRelayLimitSettings {
   2293                 max_pending_events: per_connection_outbound_queue,
   2294                 max_subscription_id_length: 64,
   2295                 max_subscriptions: 64,
   2296                 max_filters_per_request: 10,
   2297                 max_tag_values_per_filter: 100,
   2298                 max_query_complexity: 610,
   2299                 max_event_tags: 200,
   2300                 max_content_length: 65_536,
   2301                 max_limit: 500,
   2302                 default_limit: 100,
   2303             })?,
   2304             16,
   2305             per_connection_outbound_queue,
   2306         )
   2307     }
   2308 
   2309     fn temp_root(name: &str) -> PathBuf {
   2310         std::env::temp_dir().join(format!("tangle-session-{name}-{}", std::process::id()))
   2311     }
   2312 }