jazz-rs 0.1.0

A framework for CRDT based, end-to-end enrypted distributed apps
Documentation
use std::{
    cell::RefCell,
    collections::{HashMap, HashSet},
    rc::Rc,
};

use conundrum::{purpose, symm_encr::EncrKey};
use credo::{ClaimKind, Credo, GroupSecretRecipientState, GroupState};
use futures::channel::mpsc::UnboundedReceiver;
use jmbl::{
    stream_format::{OpStreamReader, SimpleOpStreamReader, SimpleOpStreamWriter, StreamReadError},
    Input, ProxyContextRef, WritableLog, JMBL,
};
use litl::{Litl, ReadError};
use jazz_telepathy::{logs::LogID, TelepathyNode};

use crate::ScopedDocID;

purpose!(LogEncryption);

mod log_reader;
mod log_writer;

use crate::conventions::{
    claim_to_declare_log_part_of_doc, claim_to_reveal_log_encr_key, log_key_secret_name,
    path_for_doc_branch,
};

pub struct ManagedJMBL {
    content_telepathy: Rc<RefCell<TelepathyNode>>,
    credo: Rc<RefCell<Credo>>,
    scoped_doc_id: ScopedDocID,
    team_update_receiver: UnboundedReceiver<GroupState>,
    current_read_logs: HashMap<LogID, SimpleOpStreamReader<log_reader::LogReader>>,
    last_secret_state: Option<GroupSecretRecipientState>,
    state: ManagedJMBLState,
    pub jmbl: JMBL,
    change_callbacks: Vec<Box<dyn FnMut()>>,
}

#[derive(Copy, Clone, PartialEq, Eq)]
enum ManagedJMBLState {
    Uninitialized,
    Loaded,
    Writable,
    Writing,
}

impl ManagedJMBL {
    pub fn id(&self) -> &ScopedDocID {
        &self.scoped_doc_id
    }

    pub fn add_change_callback<F>(&mut self, callback: F)
    where
        F: FnMut() + 'static,
    {
        self.change_callbacks.push(Box::new(callback));
    }

    pub fn load(
        scoped_doc_id: ScopedDocID,
        content_telepathy: Rc<RefCell<TelepathyNode>>,
        credo: Rc<RefCell<Credo>>,
    ) -> ManagedJMBL {
        let (team_update_listener, team_update_receiver) = futures::channel::mpsc::unbounded();
        credo
            .borrow_mut()
            .subscribe(scoped_doc_id.team, Box::new(team_update_listener));

        ManagedJMBL {
            content_telepathy,
            credo,
            scoped_doc_id,
            team_update_receiver,
            last_secret_state: None,
            current_read_logs: HashMap::new(),
            state: ManagedJMBLState::Uninitialized,
            jmbl: JMBL::new_empty(),
            change_callbacks: Vec::new()
        }
    }

    pub fn create<I: Into<Input>>(
        scoped_doc_id: ScopedDocID,
        input: I,
        content_telepathy: Rc<RefCell<TelepathyNode>>,
        credo: Rc<RefCell<Credo>>,
    ) -> ManagedJMBL {
        let mut managed = ManagedJMBL::load(scoped_doc_id, content_telepathy, credo);

        managed.receive_updates();

        assert!(managed.last_secret_state.is_some());
        assert!(managed.state == ManagedJMBLState::Uninitialized);

        managed.jmbl = JMBL::new_from_root(input, managed.make_new_log());
        managed.state = ManagedJMBLState::Writable;

        managed
    }

    pub(crate) fn make_new_log(&self) -> WritableLog {
        let log_write_access = self
            .content_telepathy
            .borrow_mut()
            .local_state
            .logs
            .create_log();
        let log_encr_key = EncrKey::new_random();
        let mut credo = self.credo.borrow_mut();
        let group_secret = credo
            .current_group_secret_for(&self.scoped_doc_id.team)
            .expect("Need to have access to group secret to create managed JMBL");

        let secret_claim_id = credo
            .make_claim_after_frontier(
                &self.scoped_doc_id.team,
                claim_to_reveal_log_encr_key(group_secret, log_write_access.id(), &log_encr_key),
            )
            .unwrap();

        credo
            .make_claim(
                &self.scoped_doc_id.team,
                claim_to_declare_log_part_of_doc(self.scoped_doc_id.clone(), log_write_access.id()),
                vec![secret_claim_id],
            )
            .unwrap();

        WritableLog::new(
            Box::new(SimpleOpStreamWriter::new(std::io::BufWriter::new(log_writer::LogWriter::new(
                Rc::clone(&self.content_telepathy),
                log_write_access,
                log_encr_key,
            )))),
        )
    }

    pub fn receive_updates(&mut self) {
        self.credo.borrow_mut().receive_from_telepathy();
        loop {
            match self.team_update_receiver.try_next() {
                Ok(Some(team_update)) => {
                    if self.last_secret_state != Some(team_update.secret_state()) {
                        self.state = match self.state {
                            ManagedJMBLState::Writing => {
                                panic!("Tried to receive updates while writing")
                            }
                            ManagedJMBLState::Writable => {
                                self.jmbl.make_readable();
                                ManagedJMBLState::Loaded
                            }
                            current => current,
                        };
                        self.last_secret_state = Some(team_update.secret_state())
                    }

                    let expected_path = path_for_doc_branch(self.scoped_doc_id.clone());

                    let logs_according_to_update = team_update
                        .valid_claims
                        .iter()
                        .filter_map(|(_id, claim)| {
                            if let ClaimKind::Statement { path, value } = &claim.expect_v1().kind {
                                if path == &expected_path {
                                    value.clone().try_into_de::<LogID>().ok()
                                } else {
                                    None
                                }
                            } else {
                                None
                            }
                        })
                        .collect::<HashSet<_>>();

                    let logs_already_followed = self
                        .current_read_logs
                        .keys()
                        .cloned()
                        .collect::<HashSet<_>>();

                    let logs_not_yet_followed =
                        logs_according_to_update.difference(&logs_already_followed);
                    let mut logs_no_longer_valid = logs_already_followed
                        .difference(&logs_according_to_update)
                        .peekable();

                    if logs_no_longer_valid.peek().is_some() {
                        unimplemented!("Reload document when logs stop being valid");
                    }

                    for log_id in logs_not_yet_followed {
                        let (listener, receiver) = futures::channel::mpsc::unbounded();
                        self.content_telepathy
                            .borrow_mut()
                            .local_state
                            .logs
                            .add_listener(*log_id, Box::new(listener));
                        let log_encr_key_litl = self
                            .credo
                            .borrow_mut()
                            .try_decrypt_entrusted_secret(
                                &self.scoped_doc_id.team,
                                &log_key_secret_name(*log_id),
                            )
                            .expect("Should be able to read log secret once new log is added");
                        self.current_read_logs.insert(
                            *log_id,
                            SimpleOpStreamReader::new(log_reader::LogReader::new(
                                receiver,
                                Litl::try_into_de(log_encr_key_litl)
                                    .expect("Expected log encryption key to deserialize"),
                            )),
                        );
                    }
                }
                Ok(None) => panic!("Team update channel closed"),
                Err(_) => {
                    break;
                }
            }
        }

        let mut received_ops = false;

        for (_log_id, log_reader) in self.current_read_logs.iter_mut() {
            loop {
                match log_reader.read_op() {
                    Ok(op) => {
                        if self.state == ManagedJMBLState::Uninitialized {
                            self.state = ManagedJMBLState::Loaded;
                        }
                        self.jmbl.apply_ops(&Some(op));
                        received_ops = true;
                    }
                    Err(err) => match err {
                        StreamReadError::LitlReadError(ReadError::Io(io_err))
                            if io_err.kind() == std::io::ErrorKind::Interrupted =>
                        {
                            break;
                        }
                        _ => panic!("Unexpected error reading log {:?}", err),
                    },
                }
            }
        }

        if received_ops {
            for callback in &mut self.change_callbacks {
                callback();
            }
        }
    }

    pub fn start_writing(&mut self) -> ProxyContextRef {
        match self.state {
            ManagedJMBLState::Writable => {}
            ManagedJMBLState::Loaded => {
                self.jmbl = self.jmbl.switch_to_new_writable_log(self.make_new_log());
                self.state = ManagedJMBLState::Writable
            }
            ManagedJMBLState::Writing => panic!("Already writing"),
            ManagedJMBLState::Uninitialized => unreachable!(),
        };

        self.state = ManagedJMBLState::Writing;
        self.jmbl.start_changing_object()
    }

    pub fn finish_writing(&mut self, write_proxy_ctx: ProxyContextRef) {
        match self.state {
            ManagedJMBLState::Writing => {
                self.jmbl.finish_changing_object(write_proxy_ctx);
                self.state = ManagedJMBLState::Writable;
                for callback in &mut self.change_callbacks {
                    callback();
                }
            }
            _ => panic!("Expected to be writing when finishing writing"),
        }
    }
}