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 }