tlpt 0.8.0

A set of protocols for building local-first distributed systems
Documentation
use std::{cell::RefCell, collections::HashMap, rc::Rc};

use audi::{Listener, ListenerSet};
use endr::ObjectID;
use futures::{channel::mpsc::channel, future, stream, FutureExt, StreamExt};
use serde_derive::{Deserialize, Serialize};

use crate::{doc::WeakDoc, ContentType};

struct BinaryStreamContentInner {
    weak_doc: Option<WeakDoc>,
    stream_listeners: ListenerSet<BinaryStreamItem>,
    mime_type: Option<String>,
    data: Vec<u8>,
    done: bool,
    last_write_log_id: Option<ObjectID>,
}

#[derive(Clone, Serialize, Deserialize, Debug)]
pub struct BinaryStreamItem {
    #[serde(skip_serializing_if = "Option::is_none")]
    pub mime_type: Option<String>,
    #[serde(with = "litl::raw_data_serde")]
    pub data: Vec<u8>,
    #[serde(skip_serializing_if = "Option::is_none")]
    pub done: Option<bool>,
}

#[derive(Clone)]
pub struct BinaryStreamContent(Rc<RefCell<BinaryStreamContentInner>>);

impl BinaryStreamContent {
    pub fn new_empty() -> BinaryStreamContent {
        BinaryStreamContent(Rc::new(RefCell::new(BinaryStreamContentInner {
            weak_doc: None,
            stream_listeners: ListenerSet::new(),
            mime_type: None,
            data: Vec::new(),
            done: false,
            last_write_log_id: None,
        })))
    }

    pub async fn add_stream_listener(&self, listener: Listener<BinaryStreamItem>) {
        let (listeners, initial_message) = {
            let self_ref = self.0.borrow();

            let initial_message = if self_ref.mime_type.is_none() {
                None
            } else {
                Some(BinaryStreamItem {
                    mime_type: Some(self_ref.mime_type.clone().unwrap()),
                    data: self_ref.data.clone(),
                    done: Some(self_ref.done),
                })
            };
            (self_ref.stream_listeners.clone(), initial_message)
        };

        listeners
            .add_with_initial_msg(listener, initial_message)
            .await;
    }

    pub fn receive_stream(
        &self,
        listener_prefix: String,
    ) -> impl futures::Stream<Item = BinaryStreamItem> {
        let (updates_tx, updates_rx) = channel(100);

        let self_clone = self.clone();

        stream::once(async move {
            self_clone
                .add_stream_listener(Listener::new(
                    &format!("{}_{:?}", listener_prefix, rand07::random::<u64>()),
                    updates_tx,
                ))
                .await;

            updates_rx
        })
        .flatten()
        .boxed_local()
    }

    pub fn start_writing(&self, mime_type: &str) {
        let mut self_ref = self.0.borrow_mut();
        if self_ref.last_write_log_id.is_some() {
            panic!("Already writing");
        }

        let (into_log, log_id) = self_ref
            .weak_doc
            .as_ref()
            .unwrap()
            .upgrade()
            .unwrap()
            .start_writing();
        self_ref.last_write_log_id = Some(log_id);

        let item = BinaryStreamItem {
            mime_type: Some(mime_type.to_owned()),
            data: vec![],
            done: None,
        };

        into_log
            .unbounded_send(litl::to_val(item).unwrap())
            .unwrap();
    }

    pub fn write(&self, data: &[u8]) {
        let self_ref = self.0.borrow_mut();

        let (into_log, log_id) = self_ref
            .weak_doc
            .as_ref()
            .unwrap()
            .upgrade()
            .unwrap()
            .start_writing();

        if Some(log_id) != self_ref.last_write_log_id {
            panic!("Log switching not supported yet for BinaryStreamContent");
        }

        let item = BinaryStreamItem {
            mime_type: None,
            data: data.to_owned(),
            done: None,
        };

        into_log
            .unbounded_send(litl::to_val(item).unwrap())
            .unwrap();
    }

    pub fn finish(&self) {
        let (into_log, log_id) = self.0.borrow_mut()
            .weak_doc
            .as_ref()
            .unwrap()
            .upgrade()
            .unwrap()
            .start_writing();

        if Some(log_id) != self.0.borrow_mut().last_write_log_id {
            panic!("Log switching not supported yet for BinaryStreamContent");
        }

        let item = BinaryStreamItem {
            mime_type: None,
            data: vec![],
            done: Some(true),
        };

        into_log
            .unbounded_send(litl::to_val(item).unwrap())
            .unwrap();
    }
}

impl ContentType for BinaryStreamContent {
    fn content_type(&self) -> &'static str {
        "binary_stream1"
    }

    fn connect_and_init(&self, weak_doc: crate::doc::WeakDoc, require_intro: bool) {
        let mut self_ref = self.0.borrow_mut();

        self_ref.weak_doc = Some(weak_doc);
    }

    fn follow_new_log(
        &self,
        log_id: endr::ObjectID,
        diffs: std::pin::Pin<Box<dyn futures::Stream<Item = super::ContentDiff>>>,
    ) -> std::pin::Pin<Box<dyn futures::Future<Output = ()>>> {
        if !self.0.borrow().data.is_empty() {
            panic!("BinaryStreamContent can only follow one log for now");
        }

        let self_for_log = self.clone();

        diffs
            .filter(|diff| future::ready(!diff.decrypted_entries.is_empty()))
            .for_each(move |diff| {
                let self_for_log = self_for_log.clone();
                async move {
                    for item_val in diff.decrypted_entries {
                        let item = litl::from_val::<BinaryStreamItem>(item_val).unwrap();
                        let listeners = {
                            let mut self_ref = self_for_log.0.borrow_mut();
                            let listeners = self_ref.stream_listeners.clone();

                            if let Some(mime_type) = &item.mime_type {
                                self_ref.mime_type = Some(mime_type.clone());
                            }
                            self_ref.data.extend(&item.data);
                            if let Some(done) = item.done {
                                self_ref.done = done;
                            }
                            listeners
                        };

                        listeners.broadcast(item).await
                    }
                }
            })
            .boxed_local()
    }
}