use std::collections::HashMap;
use std::hash::Hash;
use std::sync::{Arc, Mutex};
use crate::cell::Computed;
use crate::queue::{
TopicDurability, TopicSnapshot, TopicSubscribeOutcome, TopicSubscriptionSnapshot,
};
use crate::thread_safe::ThreadSafeContext;
use crate::topic_core::TopicCore;
struct Inner<T, I> {
core: Arc<Mutex<TopicCore<T, I>>>,
readers: Mutex<HashMap<I, Computed<Vec<T>>>>,
}
pub struct ThreadSafeTopicCell<T, I = String> {
inner: Arc<Inner<T, I>>,
}
impl<T, I> Clone for ThreadSafeTopicCell<T, I> {
fn clone(&self) -> Self {
Self {
inner: Arc::clone(&self.inner),
}
}
}
impl<T, I> ThreadSafeTopicCell<T, I>
where
T: PartialEq + Clone + Send + Sync + 'static,
I: Eq + Hash + Clone + Send + Sync + 'static,
{
pub fn new(_ctx: &ThreadSafeContext) -> Self {
Self {
inner: Arc::new(Inner {
core: Arc::new(Mutex::new(TopicCore::new())),
readers: Mutex::new(HashMap::new()),
}),
}
}
pub fn from_snapshot(ctx: &ThreadSafeContext, snapshot: TopicSnapshot<T, I>) -> Self {
let topic = Self {
inner: Arc::new(Inner {
core: Arc::new(Mutex::new(TopicCore::from_snapshot(snapshot))),
readers: Mutex::new(HashMap::new()),
}),
};
let ids = topic
.inner
.core
.lock()
.expect("topic core")
.subscription_ids();
for id in ids {
topic.ensure_reader(ctx, id);
}
topic
}
fn ensure_reader(&self, ctx: &ThreadSafeContext, id: I) -> Computed<Vec<T>> {
let mut readers = self.inner.readers.lock().expect("topic readers");
if let Some(handle) = readers.get(&id) {
return *handle;
}
let core = Arc::clone(&self.inner.core);
let reader_id = id.clone();
let handle =
ctx.computed(move |_| core.lock().expect("topic core").read_suffix(&reader_id));
readers.insert(id, handle);
handle
}
pub fn subscribe(
&self,
ctx: &ThreadSafeContext,
id: I,
durability: TopicDurability,
) -> TopicSubscribeOutcome {
let (outcome, invalidate) = {
let mut core = self.inner.core.lock().expect("topic core");
core.subscribe(id.clone(), durability)
};
let reader = self.ensure_reader(ctx, id);
if invalidate {
ctx.clear(&reader);
}
outcome
}
pub fn reconnect(&self, ctx: &ThreadSafeContext, id: I) -> TopicSubscribeOutcome {
self.subscribe(ctx, id, TopicDurability::Durable)
}
pub fn disconnect(&self, ctx: &ThreadSafeContext, id: &I) -> bool {
let remove_reader = {
let mut core = self.inner.core.lock().expect("topic core");
match core.disconnect(id) {
Some(remove_reader) => remove_reader,
None => return false,
}
};
let reader = {
let mut readers = self.inner.readers.lock().expect("topic readers");
if remove_reader {
readers.remove(id)
} else {
readers.get(id).copied()
}
};
if let Some(reader) = reader {
ctx.clear(&reader);
}
true
}
pub fn publish(&self, ctx: &ThreadSafeContext, value: T) -> u64 {
let (offset, connected) = {
let mut core = self.inner.core.lock().expect("topic core");
core.publish(value)
};
let roots: Vec<Computed<Vec<T>>> = {
let readers = self.inner.readers.lock().expect("topic readers");
connected
.iter()
.filter_map(|id| readers.get(id).copied())
.collect()
};
if !roots.is_empty() {
ctx.batch(|_| {
for reader in &roots {
ctx.clear(reader);
}
});
}
offset
}
pub fn read_stream(&self, ctx: &ThreadSafeContext, id: &I) -> Vec<T> {
let reader = self
.inner
.readers
.lock()
.expect("topic readers")
.get(id)
.copied();
reader.map(|reader| ctx.get(&reader)).unwrap_or_default()
}
pub fn read(&self, ctx: &ThreadSafeContext, id: &I) -> Option<T> {
self.read_stream(ctx, id).into_iter().next()
}
pub fn advance(&self, ctx: &ThreadSafeContext, id: &I) -> Option<T> {
let value = {
let mut core = self.inner.core.lock().expect("topic core");
core.advance(id)?
};
let reader = self
.inner
.readers
.lock()
.expect("topic readers")
.get(id)
.copied();
if let Some(reader) = reader {
ctx.clear(&reader);
}
Some(value)
}
pub fn gc(&self) -> usize {
self.inner.core.lock().expect("topic core").gc()
}
pub fn base_offset(&self) -> u64 {
self.inner.core.lock().expect("topic core").base_offset()
}
pub fn end_offset(&self) -> u64 {
self.inner.core.lock().expect("topic core").end_offset()
}
pub fn elements(&self) -> Vec<T> {
self.inner.core.lock().expect("topic core").elements()
}
pub fn subscription(&self, id: &I) -> Option<TopicSubscriptionSnapshot> {
self.inner.core.lock().expect("topic core").subscription(id)
}
pub fn reader_handle(&self, id: &I) -> Option<Computed<Vec<T>>> {
self.inner
.readers
.lock()
.expect("topic readers")
.get(id)
.copied()
}
pub fn snapshot(&self) -> TopicSnapshot<T, I> {
self.inner.core.lock().expect("topic core").snapshot()
}
}