use futures::{future, Stream, StreamExt, TryStreamExt};
use jmbl::{
ops::OpWithTarget,
stream_format::{OpDeserializer, StreamFormatReadError},
};
use litl::{ReadNewlnSepStreamError, Val};
use ridl::{symm_encr::KeySecret, unauth_symm_encr::UnauthEncryptionStream};
use thiserror::Error;
use tlpt::Diff;
use tracing::trace;
pub fn decrypting_diff_reader(
diffs: impl Stream<Item = Diff>,
log_encr_key: KeySecret,
) -> impl Stream<Item = std::io::Result<Vec<u8>>> {
let mut decryption_stream = UnauthEncryptionStream::new(log_encr_key, [0; 12].into());
diffs.map(move |diff| match diff {
Diff::Log(log_diff) => {
let mut append = log_diff.append;
decryption_stream.xor_chunk(&mut append);
Ok(append)
}
_ => panic!("Expected log diff"),
})
}
fn op_stream<S: Stream<Item = Result<Val, ReadNewlnSepStreamError>>>(
values: S,
) -> impl Stream<Item = Result<OpWithTarget, LogStreamError>> {
values.scan(OpDeserializer::new(), |deserializer, value| {
future::ready(Some(
value
.map_err(Into::into)
.and_then(|litl| deserializer.deserialize(litl).map_err(Into::into)),
))
})
}
pub fn log_stream(
receiver: impl Stream<Item = Diff>,
log_encr_key: KeySecret,
) -> impl Stream<Item = Result<OpWithTarget, LogStreamError>> {
op_stream(litl::read_newln_sep_stream(
decrypting_diff_reader(receiver, log_encr_key).into_async_read(),
))
.inspect_ok(|op| {
let age = ti64::now().0 - op.op.time.0;
trace!(age_ms = age, "Op age");
})
}
#[derive(Error, Debug)]
pub enum LogStreamError {
#[error(transparent)]
StreamFormatReadError(#[from] StreamFormatReadError),
#[error(transparent)]
ReadError(#[from] ReadNewlnSepStreamError),
}