hyperlane-plugin-websocket 21.12.0

A WebSocket plugin for the Hyperlane framework, providing robust WebSocket communication capabilities and integrating with hyperlane-broadcast for efficient message dissemination.
Documentation
use super::*;

#[test]
fn a_new_web_socket_exposes_an_empty_broadcast_map() {
    let created: WebSocket = WebSocket::new();
    let defaulted: WebSocket = WebSocket::default();
    let created_missing: Option<ReceiverCount> =
        created.get_broadcast_map().receiver_count("absent");
    let defaulted_missing: Option<ReceiverCount> =
        defaulted.get_broadcast_map().receiver_count("absent");
    assert_eq!(created_missing, None);
    assert_eq!(defaulted_missing, None);
}

#[test]
fn receiver_count_is_zero_before_anything_subscribes() {
    let web_socket: WebSocket = WebSocket::new();
    let group: ReceiverCount =
        web_socket.receiver_count(BroadcastType::<String>::PointToGroup(String::from("room")));
    let pair: ReceiverCount = web_socket.receiver_count(BroadcastType::<String>::PointToPoint(
        String::from("a"),
        String::from("b"),
    ));
    let unknown: ReceiverCount = web_socket.receiver_count(BroadcastType::<String>::Unknown);
    assert_eq!(group, 0);
    assert_eq!(pair, 0);
    assert_eq!(unknown, 0);
}

#[test]
fn receiver_count_tracks_every_subscriber_of_the_generated_key() {
    let web_socket: WebSocket = WebSocket::new();
    let broadcast_type: BroadcastType<String> = BroadcastType::PointToGroup(String::from("room"));
    let key: String = BroadcastType::get_key(broadcast_type.clone());
    let _first: BroadcastMapReceiver<Vec<u8>> = web_socket
        .get_broadcast_map()
        .subscribe_or_insert(&key, DEFAULT_BROADCAST_SENDER_CAPACITY);
    let _second: BroadcastMapReceiver<Vec<u8>> = web_socket
        .get_broadcast_map()
        .subscribe_or_insert(&key, DEFAULT_BROADCAST_SENDER_CAPACITY);
    let count: ReceiverCount = web_socket.receiver_count(broadcast_type);
    assert_eq!(count, 2);
}

#[test]
fn both_point_to_point_orders_reach_the_same_channel() {
    let web_socket: WebSocket = WebSocket::new();
    let forward: BroadcastType<String> =
        BroadcastType::PointToPoint(String::from("alice"), String::from("bob"));
    let reverse: BroadcastType<String> =
        BroadcastType::PointToPoint(String::from("bob"), String::from("alice"));
    let key: String = BroadcastType::get_key(forward.clone());
    let _receiver: BroadcastMapReceiver<Vec<u8>> = web_socket
        .get_broadcast_map()
        .subscribe_or_insert(&key, DEFAULT_BROADCAST_SENDER_CAPACITY);
    let forward_count: ReceiverCount = web_socket.receiver_count(forward);
    let reverse_count: ReceiverCount = web_socket.receiver_count(reverse);
    assert_eq!(key, String::from("ptp--alice-bob"));
    assert_eq!(forward_count, 1);
    assert_eq!(reverse_count, 1);
}

#[test]
fn receiver_count_before_connected_reports_one_more_than_the_live_count() {
    let web_socket: WebSocket = WebSocket::new();
    let broadcast_type: BroadcastType<String> = BroadcastType::PointToGroup(String::from("room"));
    let key: String = BroadcastType::get_key(broadcast_type.clone());
    let before_any: ReceiverCount =
        web_socket.receiver_count_before_connected(broadcast_type.clone());
    let _first: BroadcastMapReceiver<Vec<u8>> = web_socket
        .get_broadcast_map()
        .subscribe_or_insert(&key, DEFAULT_BROADCAST_SENDER_CAPACITY);
    let after_one: ReceiverCount =
        web_socket.receiver_count_before_connected(broadcast_type.clone());
    let _second: BroadcastMapReceiver<Vec<u8>> = web_socket
        .get_broadcast_map()
        .subscribe_or_insert(&key, DEFAULT_BROADCAST_SENDER_CAPACITY);
    let _third: BroadcastMapReceiver<Vec<u8>> = web_socket
        .get_broadcast_map()
        .subscribe_or_insert(&key, DEFAULT_BROADCAST_SENDER_CAPACITY);
    let after_three: ReceiverCount =
        web_socket.receiver_count_before_connected(broadcast_type.clone());
    let untouched: ReceiverCount =
        web_socket
            .receiver_count_before_connected(BroadcastType::<String>::PointToGroup(String::new()));
    assert_eq!(before_any, 1);
    assert_eq!(after_one, 2);
    assert_eq!(after_three, 4);
    assert_eq!(untouched, 1);
}

#[test]
fn receiver_count_after_closed_never_drops_below_zero() {
    let web_socket: WebSocket = WebSocket::new();
    let broadcast_type: BroadcastType<String> = BroadcastType::PointToGroup(String::from("room"));
    let key: String = BroadcastType::get_key(broadcast_type.clone());
    let with_no_subscriber: ReceiverCount =
        web_socket.receiver_count_after_closed(broadcast_type.clone());
    let _first: BroadcastMapReceiver<Vec<u8>> = web_socket
        .get_broadcast_map()
        .subscribe_or_insert(&key, DEFAULT_BROADCAST_SENDER_CAPACITY);
    let with_one_subscriber: ReceiverCount =
        web_socket.receiver_count_after_closed(broadcast_type.clone());
    let _second: BroadcastMapReceiver<Vec<u8>> = web_socket
        .get_broadcast_map()
        .subscribe_or_insert(&key, DEFAULT_BROADCAST_SENDER_CAPACITY);
    let with_two_subscribers: ReceiverCount =
        web_socket.receiver_count_after_closed(broadcast_type);
    let untouched: ReceiverCount = web_socket
        .receiver_count_after_closed(BroadcastType::<String>::PointToGroup(String::new()));
    assert_eq!(with_no_subscriber, 0);
    assert_eq!(with_one_subscriber, 0);
    assert_eq!(with_two_subscribers, 1);
    assert_eq!(untouched, 0);
}

#[test]
fn try_send_reports_no_channel_for_an_unused_broadcast_type() {
    let web_socket: WebSocket = WebSocket::new();
    let result: Result<Option<ReceiverCount>, BroadcastMapSendError<Vec<u8>>> = web_socket
        .try_send(
            BroadcastType::<String>::PointToGroup(String::from("room")),
            "payload",
        );
    let count: Option<ReceiverCount> = result.unwrap();
    assert_eq!(count, None);
}

#[tokio::test]
async fn try_send_delivers_the_payload_to_every_subscriber() {
    let web_socket: WebSocket = WebSocket::new();
    let broadcast_type: BroadcastType<String> = BroadcastType::PointToGroup(String::from("room"));
    let key: String = BroadcastType::get_key(broadcast_type.clone());
    let mut first: BroadcastMapReceiver<Vec<u8>> = web_socket
        .get_broadcast_map()
        .subscribe_or_insert(&key, DEFAULT_BROADCAST_SENDER_CAPACITY);
    let mut second: BroadcastMapReceiver<Vec<u8>> = web_socket
        .get_broadcast_map()
        .subscribe_or_insert(&key, DEFAULT_BROADCAST_SENDER_CAPACITY);
    let result: Result<Option<ReceiverCount>, BroadcastMapSendError<Vec<u8>>> =
        web_socket.try_send(broadcast_type, String::from("hello"));
    let count: Option<ReceiverCount> = result.unwrap();
    assert_eq!(count, Some(2));
    let first_message: Result<Vec<u8>, RecvError> = first.recv().await;
    let second_message: Result<Vec<u8>, RecvError> = second.recv().await;
    assert_eq!(first_message.unwrap(), String::from("hello").into_bytes());
    assert_eq!(second_message.unwrap(), String::from("hello").into_bytes());
}

#[tokio::test]
async fn try_send_accepts_every_byte_vector_payload_shape() {
    let web_socket: WebSocket = WebSocket::new();
    let broadcast_type: BroadcastType<String> = BroadcastType::PointToGroup(String::from("room"));
    let key: String = BroadcastType::get_key(broadcast_type.clone());
    let mut receiver: BroadcastMapReceiver<Vec<u8>> = web_socket
        .get_broadcast_map()
        .subscribe_or_insert(&key, DEFAULT_BROADCAST_SENDER_CAPACITY);
    let from_str: Result<Option<ReceiverCount>, BroadcastMapSendError<Vec<u8>>> =
        web_socket.try_send(broadcast_type.clone(), "one");
    let from_string: Result<Option<ReceiverCount>, BroadcastMapSendError<Vec<u8>>> =
        web_socket.try_send(broadcast_type.clone(), String::from("two"));
    let from_bytes: Result<Option<ReceiverCount>, BroadcastMapSendError<Vec<u8>>> =
        web_socket.try_send(broadcast_type, String::from("three").into_bytes());
    assert_eq!(from_str.unwrap(), Some(1));
    assert_eq!(from_string.unwrap(), Some(1));
    assert_eq!(from_bytes.unwrap(), Some(1));
    let first: Result<Vec<u8>, RecvError> = receiver.recv().await;
    let second: Result<Vec<u8>, RecvError> = receiver.recv().await;
    let third: Result<Vec<u8>, RecvError> = receiver.recv().await;
    assert_eq!(first.unwrap(), String::from("one").into_bytes());
    assert_eq!(second.unwrap(), String::from("two").into_bytes());
    assert_eq!(third.unwrap(), String::from("three").into_bytes());
}

#[test]
fn try_send_fails_and_hands_the_payload_back_when_no_subscriber_is_left() {
    let web_socket: WebSocket = WebSocket::new();
    let broadcast_type: BroadcastType<String> = BroadcastType::PointToGroup(String::from("room"));
    let key: String = BroadcastType::get_key(broadcast_type.clone());
    let receiver: BroadcastMapReceiver<Vec<u8>> = web_socket
        .get_broadcast_map()
        .subscribe_or_insert(&key, DEFAULT_BROADCAST_SENDER_CAPACITY);
    drop(receiver);
    let result: Result<Option<ReceiverCount>, BroadcastMapSendError<Vec<u8>>> =
        web_socket.try_send(broadcast_type, String::from("hello"));
    let error: BroadcastMapSendError<Vec<u8>> = result.unwrap_err();
    let message: String = error.to_string();
    assert_eq!(message, String::from("channel closed"));
    let recovered: Vec<u8> = error.0;
    assert_eq!(recovered, String::from("hello").into_bytes());
}

#[test]
fn send_returns_the_number_of_receivers() {
    let web_socket: WebSocket = WebSocket::new();
    let broadcast_type: BroadcastType<String> = BroadcastType::PointToGroup(String::from("room"));
    let key: String = BroadcastType::get_key(broadcast_type.clone());
    let _receiver: BroadcastMapReceiver<Vec<u8>> = web_socket
        .get_broadcast_map()
        .subscribe_or_insert(&key, DEFAULT_BROADCAST_SENDER_CAPACITY);
    let count: Option<ReceiverCount> = web_socket.send(broadcast_type, "hello");
    assert_eq!(count, Some(1));
}

#[test]
fn send_returns_none_for_a_channel_that_was_never_created() {
    let web_socket: WebSocket = WebSocket::new();
    let count: Option<ReceiverCount> = web_socket.send(
        BroadcastType::<String>::PointToGroup(String::from("room")),
        "hello",
    );
    assert_eq!(count, None);
}

#[test]
#[should_panic(expected = "SendError")]
fn send_panics_when_the_channel_has_no_subscriber() {
    let web_socket: WebSocket = WebSocket::new();
    let broadcast_type: BroadcastType<String> = BroadcastType::PointToGroup(String::from("room"));
    let key: String = BroadcastType::get_key(broadcast_type.clone());
    let receiver: BroadcastMapReceiver<Vec<u8>> = web_socket
        .get_broadcast_map()
        .subscribe_or_insert(&key, DEFAULT_BROADCAST_SENDER_CAPACITY);
    drop(receiver);
    let count: Option<ReceiverCount> = web_socket.send(broadcast_type, "hello");
    assert_eq!(count, None);
}

#[test]
fn unsubscribing_a_key_removes_the_whole_channel() {
    let web_socket: WebSocket = WebSocket::new();
    let broadcast_type: BroadcastType<String> = BroadcastType::PointToGroup(String::from("room"));
    let key: String = BroadcastType::get_key(broadcast_type.clone());
    let _receiver: BroadcastMapReceiver<Vec<u8>> = web_socket
        .get_broadcast_map()
        .subscribe_or_insert(&key, DEFAULT_BROADCAST_SENDER_CAPACITY);
    let removed: Option<Broadcast<Vec<u8>>> = web_socket.get_broadcast_map().unsubscribe(&key);
    assert!(removed.is_some());
    let count: ReceiverCount = web_socket.receiver_count(broadcast_type.clone());
    assert_eq!(count, 0);
    let result: Result<Option<ReceiverCount>, BroadcastMapSendError<Vec<u8>>> =
        web_socket.try_send(broadcast_type, "hello");
    let sent: Option<ReceiverCount> = result.unwrap();
    assert_eq!(sent, None);
    let removed_again: Option<Broadcast<Vec<u8>>> =
        web_socket.get_broadcast_map().unsubscribe(&key);
    assert!(removed_again.is_none());
}

#[test]
fn a_cloned_web_socket_keeps_counting_the_subscribers_of_the_original() {
    let web_socket: WebSocket = WebSocket::new();
    let broadcast_type: BroadcastType<String> = BroadcastType::PointToGroup(String::from("room"));
    let key: String = BroadcastType::get_key(broadcast_type.clone());
    let _receiver: BroadcastMapReceiver<Vec<u8>> = web_socket
        .get_broadcast_map()
        .subscribe_or_insert(&key, DEFAULT_BROADCAST_SENDER_CAPACITY);
    let clone: WebSocket = web_socket.clone();
    let count: ReceiverCount = clone.receiver_count(broadcast_type);
    assert_eq!(count, 1);
}

#[test]
fn a_cloned_web_socket_does_not_share_later_map_insertions() {
    let web_socket: WebSocket = WebSocket::new();
    let clone: WebSocket = web_socket.clone();
    let broadcast_type: BroadcastType<String> = BroadcastType::PointToGroup(String::from("room"));
    let key: String = BroadcastType::get_key(broadcast_type.clone());
    let _receiver: BroadcastMapReceiver<Vec<u8>> = clone
        .get_broadcast_map()
        .subscribe_or_insert(&key, DEFAULT_BROADCAST_SENDER_CAPACITY);
    let clone_count: ReceiverCount = clone.receiver_count(broadcast_type.clone());
    let original_count: ReceiverCount = web_socket.receiver_count(broadcast_type);
    assert_eq!(clone_count, 1);
    assert_eq!(original_count, 0);
}

#[test]
fn debug_renders_the_inner_broadcast_map_without_panicking() {
    let web_socket: WebSocket = WebSocket::new();
    let empty: String = format!("{:?}", web_socket);
    assert!(!empty.is_empty());
    let key: String =
        BroadcastType::<String>::get_key(BroadcastType::PointToGroup(String::from("room")));
    let _receiver: BroadcastMapReceiver<Vec<u8>> = web_socket
        .get_broadcast_map()
        .subscribe_or_insert(&key, DEFAULT_BROADCAST_SENDER_CAPACITY);
    let populated: String = format!("{:?}", web_socket);
    assert!(populated.contains("ptg--room"));
    assert_ne!(empty, populated);
}