jazz-rs 0.2.8

A framework for CRDT based, end-to-end enrypted distributed apps
Documentation
use futures::{
    channel::mpsc::Receiver, future, stream, AsyncRead, Stream, StreamExt, TryStreamExt,
};
use jmbl::{
    ops::OpWithTarget,
    stream_format::{OpDeserializer, StreamReadError},
};
use litl::{Litl, ReadError};
use ridl::{symm_encr::KeySecret, unauth_symm_encr::UnauthEncryptionStream};
use tlpt::Diff;

struct DecryptingDiffReader {
    receiver: Receiver<Diff>,
    decryption_stream: UnauthEncryptionStream,
}

impl DecryptingDiffReader {
    pub fn new(receiver: Receiver<Diff>, log_encr_key: KeySecret) -> Self {
        Self {
            receiver,
            // Null nonce is safe IFF the key is unique to this log
            decryption_stream: UnauthEncryptionStream::new(log_encr_key, [0; 12].into()),
        }
    }
}

// this should make it possible to use the DecryptingReader as an async reader
impl Stream for DecryptingDiffReader {
    type Item = Result<Vec<u8>, futures::io::Error>;

    fn poll_next(
        mut self: std::pin::Pin<&mut Self>,
        cx: &mut std::task::Context<'_>,
    ) -> std::task::Poll<Option<Self::Item>> {
        match self.receiver.poll_next_unpin(cx) {
            std::task::Poll::Ready(Some(diff)) => match diff {
                Diff::Log(log_diff) => {
                    let mut append = log_diff.append;
                    self.decryption_stream.xor_chunk(&mut append);
                    std::task::Poll::Ready(Some(Ok(append)))
                }
                _ => panic!("Expected log diff"),
            },
            std::task::Poll::Ready(None) => std::task::Poll::Ready(None),
            std::task::Poll::Pending => std::task::Poll::Pending,
        }
    }
}

fn op_stream<S: Stream<Item = Result<Litl, ReadError>>>(
    litls: S,
) -> impl Stream<Item = Result<OpWithTarget, StreamReadError>> {
    litls.scan(OpDeserializer::new(), |deserializer, litl| {
        future::ready(Some(
            litl.map_err(|read_err| read_err.into())
                .and_then(|litl| deserializer.deserialize(litl)),
        ))
    })
}

pub fn log_stream(
    receiver: Receiver<Diff>,
    log_encr_key: KeySecret,
) -> impl Stream<Item = Result<OpWithTarget, StreamReadError>> {
    // TODO: add decompression
    op_stream(Litl::read_stream(
        DecryptingDiffReader::new(receiver, log_encr_key).into_async_read(),
    ))
}