rings-node 0.20.0

Rings is a structured peer-to-peer network implementation using WebRTC, Chord algorithm, and full WebAssembly (WASM) support.
Documentation
use std::sync::Arc;

use async_trait::async_trait;
use futures::lock::Mutex;
use rings_core::dht::Did;
use rings_core::inspect::SwarmInspect;
use rings_core::swarm::callback::SwarmCallback;
use rings_core::utils;
use wasm_bindgen_futures::spawn_local;
use wasm_bindgen_test::*;

use crate::prelude::*;
use crate::processor::*;
use crate::tests::wasm::prepare_processor;

async fn close_all_connections(p: &Processor) {
    futures::future::join_all(
        p.swarm
            .peers()
            .iter()
            .map(|peer| p.swarm.disconnect(peer.did.parse().unwrap())),
    )
    .await;
}

struct SwarmCallbackStruct {
    pub msgs: Mutex<Vec<String>>,
}

#[async_trait(?Send)]
impl SwarmCallback for SwarmCallbackStruct {
    async fn on_inbound(
        &self,
        payload: &MessagePayload,
    ) -> Result<(), rings_core::error::CallbackError> {
        let msg: Message = payload.transaction.data().map_err(Box::new)?;
        if let Message::CustomMessage(ref msg) = msg {
            let text = String::from_utf8(msg.0.to_vec()).unwrap();
            console_log!("msg received: {}", text);
            let mut msgs = self.msgs.try_lock().unwrap();
            msgs.push(text);
        }
        Ok(())
    }
}

async fn wait_for_swarm_state(
    p: &Processor,
    label: &str,
    predicate: impl Fn(&SwarmInspect) -> bool,
) {
    let mut last_inspect = None;
    for _ in 0..100 {
        let inspect = p.swarm.inspect().await;
        if predicate(&inspect) {
            return;
        }
        last_inspect = Some(inspect);
        utils::js_utils::window_sleep(200).await.unwrap();
    }
    panic!(
        "timeout waiting for {label}; peers={:?}, dht={:?}",
        last_inspect.as_ref().map(|inspect| &inspect.peers),
        last_inspect.as_ref().map(|inspect| &inspect.dht),
    );
}

/// Wait until `did` shows up in `p`'s DHT successor list.
///
/// The WebRTC data channel opens asynchronously after `accept_answer`, and the
/// DHT successors are only populated in the `on_data_channel_open -> join_dht`
/// callback. Polling the DHT here (instead of relying on a fixed sleep) makes
/// the test robust against slow connection setup on CI, where an empty
/// successor list would cause `find_successor` to route a message to the node
/// itself and fail with `SwarmMissDidInTable`.
async fn wait_for_dht_successor(p: &Processor, did: Did) {
    let did = did.to_string();
    let label = format!("{did} to appear in DHT successors");
    wait_for_swarm_state(p, &label, |inspect| {
        inspect
            .dht
            .successors
            .iter()
            .any(|successor| successor == &did)
    })
    .await;
}

async fn wait_for_single_connected_peer(p: &Processor, did: Did) {
    let did = did.to_string();
    let label = format!("{did} to have exactly one Connected peer entry");
    wait_for_swarm_state(p, &label, |inspect| {
        let mut matching_peers = inspect.peers.iter().filter(|peer| peer.did == did);
        matches!(
            (matching_peers.next(), matching_peers.next()),
            (Some(peer), None) if peer.state == "Connected"
        )
    })
    .await;
}

async fn create_connection(p1: &Processor, p2: &Processor) {
    console_log!("create_offer");
    let offer = p1.swarm.create_offer(p2.did()).await.unwrap();

    console_log!("answer_offer");
    let answer = p2.swarm.answer_offer(offer).await.unwrap();

    console_log!("accept_answer");
    p1.swarm.accept_answer(answer).await.unwrap();

    // Wait for the data channel to open and the DHT to converge on both sides
    // before returning, so callers can rely on the routing tables being ready.
    futures::join!(
        wait_for_dht_successor(p1, p2.did()),
        wait_for_dht_successor(p2, p1.did()),
    );
}

#[wasm_bindgen_test]
async fn test_processor_handshake_and_msg() {
    let callback1 = Arc::new(SwarmCallbackStruct {
        msgs: Mutex::new(vec![]),
    });
    let callback2 = Arc::new(SwarmCallbackStruct {
        msgs: Mutex::new(vec![]),
    });

    let p1 = prepare_processor().await;
    let p2 = prepare_processor().await;

    p1.swarm.set_callback(callback1.clone()).unwrap();
    p2.swarm.set_callback(callback2.clone()).unwrap();

    let test_text1 = "test1";
    let test_text2 = "test2";
    let test_text3 = "test3";
    let test_text4 = "test4";
    let test_text5 = "test5";

    let p1_did = p1.did();
    let p2_did = p2.did();
    console_log!("p1_did: {}", p1_did);
    console_log!("p2_did: {}", p2_did);

    console_log!("listen");
    let p1_listener = p1.clone();
    spawn_local(async move {
        p1_listener.listen().await;
    });
    let p2_listener = p2.clone();
    spawn_local(async move {
        p2_listener.listen().await;
    });

    console_log!("processor_hs_connect_1_2");
    create_connection(&p1, &p2).await;

    utils::js_utils::window_sleep(2000).await.unwrap();

    console_log!("processor_send_test_text_messages");
    p1.send_message(p2_did, test_text1.as_bytes())
        .await
        .unwrap();
    console_log!("send test_text1 done");

    p2.send_message(p1_did, test_text2.as_bytes())
        .await
        .unwrap();
    console_log!("send test_text2 done");

    p2.send_message(p1_did, test_text3.as_bytes())
        .await
        .unwrap();
    console_log!("send test_text3 done");

    p1.send_message(p2_did, test_text4.as_bytes())
        .await
        .unwrap();
    console_log!("send test_text4 done");

    p2.send_message(p1_did, test_text5.as_bytes())
        .await
        .unwrap();
    console_log!("send test_text5 done");

    utils::js_utils::window_sleep(4000).await.unwrap();

    console_log!("check received");

    let mut msgs1 = callback1.msgs.try_lock().unwrap().as_slice().to_vec();
    msgs1.sort();
    let mut msgs2 = callback2.msgs.try_lock().unwrap().as_slice().to_vec();
    msgs2.sort();

    let mut expect1 = vec![
        test_text2.to_owned(),
        test_text3.to_owned(),
        test_text5.to_owned(),
    ];
    expect1.sort();

    let mut expect2 = vec![test_text1.to_owned(), test_text4.to_owned()];
    expect2.sort();
    assert_eq!(msgs1, expect1);
    assert_eq!(msgs2, expect2);

    console_log!("processor_hs_close_all_connections");
    futures::join!(close_all_connections(&p1), close_all_connections(&p2),);
}

#[wasm_bindgen_test]
async fn test_processor_connect_with_did() {
    super::setup_log();
    let p1 = prepare_processor().await;
    console_log!("p1 address: {}", p1.did());
    let p2 = prepare_processor().await;
    console_log!("p2 address: {}", p2.did());
    let p3 = prepare_processor().await;
    console_log!("p3 address: {}", p3.did());

    console_log!("processor_connect_p1_and_p2");
    create_connection(&p1, &p2).await;
    console_log!("processor_connect_p1_and_p2, done");

    console_log!("processor_connect_p2_and_p3");
    create_connection(&p2, &p3).await;
    console_log!("processor_connect_p2_and_p3, done");

    let peers = p1.swarm.peers();
    assert!(
        peers.iter().any(|peer| peer.did.eq(&p2.did().to_string())),
        "p2 not in p1's peer list"
    );

    // No sleep needed here: `create_connection` already waited for both links to
    // converge in the DHT, so the relay node p2 knows p1 and p3 and can forward
    // p1's offer to p3 (and route p3's answer back to p1).

    console_log!("connect p1 and p3");
    let (p1_connect, p3_connect) =
        futures::join!(p1.connect_with_did(p3.did()), p3.connect_with_did(p1.did()),);
    p1_connect.unwrap();
    p3_connect.unwrap();
    futures::join!(
        wait_for_single_connected_peer(&p1, p3.did()),
        wait_for_single_connected_peer(&p3, p1.did()),
    );
    console_log!("processor_detect_connection_state");

    console_log!("check peers");
    let peers = p1.swarm.peers();
    assert!(
        peers
            .iter()
            .any(|peer| peer.did.eq(p3.did().to_string().as_str())),
        "peer list dose NOT contains p3 address"
    );
    console_log!("processor_close_all_connections");
    futures::join!(
        close_all_connections(&p1),
        close_all_connections(&p2),
        close_all_connections(&p3),
    );
}