use std::{
cell::RefCell,
collections::{HashMap, HashSet},
rc::Rc,
};
use audi::{Listener, ListenerSet};
use credo::{ClaimBody, Credo, Scope, ScopeSecretRecipientState, ScopeState};
use futures::{
channel::{mpsc::channel, oneshot},
stream, FutureExt, Stream, StreamExt,
};
use jmbl::{ops::OpWithTarget, Input, JMBLViewRef, JMBL};
use mofo::Mofo;
use caro::ObjectID;
use tracing::trace;
use crate::{
conventions::{log_key_secret_name, INCLUDE_LOG_STR, READ_CONTENT},
managed_jmbl::writing::create_log_as_writable_log,
};
mod reading;
mod writing;
use self::reading::log_stream;
struct ManagedJMBLInner {
doc_scope: Scope,
content: caro::Node,
last_secret_state: Option<ScopeSecretRecipientState>,
tx_first_secret_state_set: Option<oneshot::Sender<()>>,
logs_already_followed: HashSet<ObjectID>,
jmbl_logs_to_tlpt_logs: HashMap<jmbl::LogID, caro::ObjectID>,
state: ManagedJMBLState,
pub jmbl: JMBL,
change_callbacks: Vec<Box<dyn FnMut(jmbl::Value, bool)>>,
listeners: ListenerSet<(jmbl::Value, bool)>,
background: Mofo,
}
#[derive(Clone)]
pub struct ManagedJMBL(Rc<RefCell<ManagedJMBLInner>>);
#[derive(Copy, Clone, PartialEq, Eq)]
enum ManagedJMBLState {
Uninitialized,
Loaded,
Writable,
Writing,
}
impl ManagedJMBLInner {
fn receive_ops(&mut self, ops: &[OpWithTarget], from_tlpt_log: ObjectID) {
if self.state == ManagedJMBLState::Uninitialized {
self.state = ManagedJMBLState::Loaded;
}
self.jmbl.apply_ops(ops);
for op in ops {
self.jmbl_logs_to_tlpt_logs
.entry(op.op.id.log_id)
.or_insert(from_tlpt_log);
}
for callback in &mut self.change_callbacks {
callback(self.jmbl.get_root(), true);
}
}
}
impl ManagedJMBL {
pub fn scope(&self) -> Scope {
self.0.borrow().doc_scope.clone()
}
pub fn current_root(&self) -> jmbl::Value {
self.0.borrow().jmbl.get_root()
}
pub fn add_change_callback<F>(&mut self, callback: F)
where
F: FnMut(jmbl::Value, bool) + 'static,
{
(*self.0)
.borrow_mut()
.change_callbacks
.push(Box::new(callback));
}
async fn add_listener(&self, listener: Listener<(jmbl::Value, bool)>) {
let listeners = (*self.0).borrow().listeners.clone();
listeners
.add_with_initial_msg(listener, Some((self.current_root(), false)))
.await;
}
pub fn updates(
&self,
listener_prefix: String,
) -> impl Stream<Item = (jmbl::Value, bool)> + 'static {
let self_rc = self.clone();
stream::once(async move {
let (tx, rx) = channel(100);
let listener = Listener::new(
&format!("{}_{:?}", listener_prefix, rand07::random::<u64>()),
tx,
);
self_rc.add_listener(listener).await;
rx
})
.flatten()
.boxed_local()
}
pub fn load(
doc_scope: Scope,
content: caro::Node,
background: Mofo,
) -> (ManagedJMBL, oneshot::Receiver<()>) {
let (tx_first_secret_state_set, rx_first_secret_state_set) = oneshot::channel();
let managed_jmbl_rc = ManagedJMBL(Rc::new(RefCell::new(ManagedJMBLInner {
content,
doc_scope: doc_scope.clone(),
last_secret_state: None,
tx_first_secret_state_set: Some(tx_first_secret_state_set),
logs_already_followed: HashSet::new(),
jmbl_logs_to_tlpt_logs: HashMap::new(),
state: ManagedJMBLState::Uninitialized,
jmbl: JMBL::new_empty(),
change_callbacks: Vec::new(),
listeners: ListenerSet::new(),
background: background.clone(),
})));
background.add_background_task(Box::pin({
let managed_jmbl_rc = managed_jmbl_rc.clone();
let doc_updates = doc_scope.updates("managed_jmbl".to_string());
doc_updates.for_each(move |doc_update| {
let managed_jmbl_rc = managed_jmbl_rc.clone();
async move {
managed_jmbl_rc.receive_doc_update(doc_update).await;
}
})
}));
(managed_jmbl_rc, rx_first_secret_state_set)
}
pub async fn create<I: Into<Input>>(
doc_scope: Scope,
input: I,
content: caro::Node,
credo: Credo,
background: Mofo,
) -> ManagedJMBL {
let (managed_jmbl, first_secret_state_set) =
ManagedJMBL::load(doc_scope, content, background);
first_secret_state_set.await.unwrap();
{
let self_for_inserting_log_mapping = managed_jmbl.clone();
let mut managed_jmbl_ref = (*managed_jmbl.0).borrow_mut();
assert!(managed_jmbl_ref.last_secret_state.is_some());
assert!(managed_jmbl_ref.state == ManagedJMBLState::Uninitialized);
let (writable_log, writing) = create_log_as_writable_log(
managed_jmbl_ref.doc_scope.clone(),
managed_jmbl_ref.content.clone(),
self_for_inserting_log_mapping,
);
managed_jmbl_ref
.background
.add_background_task(writing.boxed_local());
managed_jmbl_ref.jmbl = JMBL::new_from_root(input, writable_log);
managed_jmbl_ref.state = ManagedJMBLState::Writable;
}
managed_jmbl
}
pub async fn receive_doc_update(&self, doc_update: ScopeState) {
let logs_not_yet_followed = {
let mut self_ref = (*self.0).borrow_mut();
if self_ref.last_secret_state != Some(doc_update.secret_state()) {
self_ref.state = match self_ref.state {
ManagedJMBLState::Writing => {
panic!("Tried to receive updates while writing")
}
ManagedJMBLState::Writable => {
self_ref.jmbl.make_readable();
ManagedJMBLState::Loaded
}
current => current,
};
if self_ref.last_secret_state.is_none() {
self_ref
.tx_first_secret_state_set
.take()
.unwrap()
.send(())
.unwrap();
}
self_ref.last_secret_state = Some(doc_update.secret_state())
}
let logs_according_to_update = doc_update
.valid_claims
.iter()
.filter_map(|(_id, (claim, _))| {
if let ClaimBody::Statement { path, value } = &claim.body {
if path == INCLUDE_LOG_STR {
litl::from_val::<ObjectID>(value.clone()).ok()
} else {
None
}
} else {
None
}
})
.collect::<HashSet<_>>();
let logs_already_followed = self_ref.logs_already_followed.clone();
let logs_not_yet_followed = logs_according_to_update
.difference(&logs_already_followed)
.cloned()
.collect::<Vec<_>>();
let mut logs_no_longer_valid = logs_already_followed
.difference(&logs_according_to_update)
.peekable();
trace!(doc_update = ?doc_update, logs_according_to_update = ?logs_according_to_update, logs_already_followed = ?logs_already_followed, logs_not_yet_followed = ?logs_not_yet_followed, "Recevied new set of document logs");
if logs_no_longer_valid.peek().is_some() {
unimplemented!("Reload document when logs stop being valid");
}
logs_not_yet_followed
};
for log_id in logs_not_yet_followed {
let content = (*self.0).borrow().content.clone();
let doc_scope = (*self.0).borrow().doc_scope.clone();
let log_rx = content.diffs(log_id, format!("managed_jmbl_reader_{}", doc_scope.id()));
let mut self_ref = (*self.0).borrow_mut();
let log_encr_key = litl::from_val(
doc_scope
.try_decrypt_entrusted_secret(READ_CONTENT, &log_key_secret_name(log_id))
.await
.expect("Should be able to read log secret once new log is added"),
)
.expect("Should be able to deserialize encryption key");
self_ref.background.add_background_task(Box::pin({
let self_rc = self.clone();
log_stream(log_rx, log_encr_key)
.map(|op| op.unwrap())
.ready_chunks(10000)
.for_each(move |ops| {
trace!(n_ops = ops.len(), "Received ops");
(*self_rc.0).borrow_mut().receive_ops(&ops, log_id);
let root = self_rc.current_root();
let listeners = (*self_rc.0).borrow().listeners.clone();
async move {
listeners.broadcast((root, true)).await;
}
})
}));
self_ref.logs_already_followed.insert(log_id);
}
}
pub fn start_writing(&self) -> JMBLViewRef {
let self_for_inserting_log_mapping = self.clone();
let mut self_ref = (*self.0).borrow_mut();
match self_ref.state {
ManagedJMBLState::Writable => {}
ManagedJMBLState::Loaded => {
let (writable_log, writing) = create_log_as_writable_log(
self_ref.doc_scope.clone(),
self_ref.content.clone(),
self_for_inserting_log_mapping,
);
self_ref
.background
.add_background_task(writing.boxed_local());
self_ref.jmbl = self_ref.jmbl.switch_to_new_writable_log(writable_log);
self_ref.state = ManagedJMBLState::Writable
}
ManagedJMBLState::Writing => panic!("Already writing"),
ManagedJMBLState::Uninitialized => unreachable!(),
};
self_ref.state = ManagedJMBLState::Writing;
self_ref.jmbl.start_changing_object()
}
pub fn finish_writing(&self, write_view: JMBLViewRef) {
let mut self_ref = (*self.0).borrow_mut();
match self_ref.state {
ManagedJMBLState::Writing => {
self_ref.jmbl.finish_changing_object(write_view);
self_ref.state = ManagedJMBLState::Writable;
let root = self_ref.jmbl.get_root();
for callback in &mut self_ref.change_callbacks {
callback(root.clone(), false);
}
}
_ => panic!("Expected to be writing when finishing writing"),
}
}
pub fn get_tlpt_log_for_jmbl_log(&self, jmbl_log: jmbl::LogID) -> Option<ObjectID> {
(*self.0)
.borrow()
.jmbl_logs_to_tlpt_logs
.get(&jmbl_log)
.cloned()
}
}