jazz-rs 0.4.0

A framework for CRDT based, end-to-end enrypted distributed apps
Documentation
use std::time::Duration;

use audi::Listener;
use futures::{Future, StreamExt};
use jazz_rs::{DocID, InvitationKind, Jazz};
use jmbl::Input;
use litl::Litl;
use mofo::Mofo;
use tlpt::Remote;
use tracing::{debug, info, trace};

#[tokio::test]
async fn can_create_team_and_share_document_for_reading() {
    traceful::init_test_trace();

    let (_, _, background, setup_and_basic_assertions) =
        create_team_and_share_document(InvitationKind::Reader);

    background.run_until(setup_and_basic_assertions).await;
}

fn create_team_and_share_document(
    invitation_kind: InvitationKind,
) -> (Jazz, Jazz, Mofo, impl Future<Output = DocID>) {
    let background = Mofo::new();
    let jazz1 = Jazz::new("jazz1".to_string(), background.clone());
    let jazz2 = Jazz::new("jazz2".to_string(), background.clone());

    let setup = {
        let jazz1 = jazz1.clone();
        let jazz2 = jazz2.clone();

        async move {
            let team = jazz1.create_team().await;
            let (doc_id, managed_doc) = jazz1
                .create_document(
                    team,
                    Input::CollabMap(vec![(
                        "test".to_string(),
                        Litl::string("Hello World").into(),
                    )]),
                )
                .await;
            assert_eq!(
                managed_doc.current_root().if_map().unwrap().to_litl(),
                Litl::dict([("test", Litl::string("Hello World"))])
            );
            let invitation_token = jazz1
                .create_invitation(team, invitation_kind, None, None, None)
                .await
                .unwrap();

            let (jazz1_as_remote, jazz2_as_remote) =
                Remote::new_connected_test_pair("jazz1", "jazz2");

            jazz1.add_remote(jazz2_as_remote).await;
            jazz2.add_remote(jazz1_as_remote).await;

            let jazz2 = jazz2.clone();

            debug!("invitation token: {:?}", invitation_token);
            jazz2.join_team(invitation_token).await.unwrap();
            let managed_doc2 = jazz2.load_document(doc_id.clone()).await;

            let (doc2_tx, doc2_rx) = futures::channel::mpsc::channel(100);

            managed_doc2
                .add_listener(Listener::new("test2", doc2_tx))
                .await;

            let (first_root, _) = doc2_rx
                .filter(|(root, _)| {
                    trace!(root = ?root, is_plain = matches!(root.if_plain(), Ok(Litl::Null)), "Got root");
                    futures::future::ready(!matches!(root.if_plain(), Ok(Litl::Null)))
                })
                .next()
                .await
                .unwrap();

            assert_eq!(
                first_root.if_map().unwrap().to_litl(),
                Litl::dict([("test", Litl::string("Hello World"))])
            );

            doc_id
        }
    };

    (jazz1, jazz2, background, setup)
}

#[tokio::test]
async fn can_create_team_and_share_document_for_writing() {
    traceful::init_test_trace();

    let (jazz1, jazz2, background, setup_and_basic_assertions) =
        create_team_and_share_document(InvitationKind::Writer);

    background
        .run_until(async {
            // TODO: #70 wait until we definitely have write access, and make waiting for that easy

            let doc_id = setup_and_basic_assertions.await;

            let managed_doc2 = jazz2.load_document(doc_id.clone()).await;

            let write_view = managed_doc2.start_writing();
            write_view
                .get_root()
                .if_map_mut()
                .unwrap()
                .insert("test2", "Litl World");
            managed_doc2.finish_writing(write_view);

            assert_eq!(
                managed_doc2.current_root().if_map().unwrap().to_litl(),
                Litl::dict([
                    ("test", Litl::string("Hello World")),
                    ("test2", Litl::string("Litl World"))
                ])
            );

            let managed_doc1 = jazz1.load_document(doc_id).await;

            let (doc1_tx, doc1_rx) = futures::channel::mpsc::channel(100);

            managed_doc1
                .add_listener(Listener::new("test1_after_sharing", doc1_tx))
                .await;

            let (first_root_with_jazz2_update, _) = doc1_rx
                .filter(|(root, _)| {
                    let has_update = root.if_map().unwrap().to_litl().get("test2").is_some();
                    trace!(root = ?root, has_update = has_update, "Got root");
                    futures::future::ready(has_update)
                })
                .next()
                .await
                .unwrap();

            assert_eq!(
                first_root_with_jazz2_update.if_map().unwrap().to_litl(),
                Litl::dict([
                    ("test", Litl::string("Hello World")),
                    ("test2", Litl::string("Litl World"))
                ])
            );
        })
        .await
}