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()
}
}