use std::collections::{HashMap, HashSet};
use std::sync::Arc;
use crate::error::Error;
use crate::topic::backend::{InMemoryBackend, TopicBackend, TopicState};
use crate::topic::handle::TopicHandle;
use crate::topic::types::TopicOptions;
use otel_arrow_dfe_config::TopicName;
use parking_lot::RwLock;
#[derive(Clone)]
pub struct TopicBroker<T: Send + Sync + 'static> {
inner: Arc<BrokerInner<T>>,
}
struct BrokerInner<T: Send + Sync + 'static> {
topics: RwLock<HashMap<TopicName, Arc<dyn TopicState<T>>>>,
}
impl<T: Send + Sync + 'static> TopicBroker<T> {
#[must_use]
pub fn new() -> Self {
Self {
inner: Arc::new(BrokerInner {
topics: RwLock::new(HashMap::new()),
}),
}
}
pub fn create_topic(
&self,
name: impl Into<TopicName>,
opts: TopicOptions,
backend: impl TopicBackend<T>,
) -> Result<TopicHandle<T>, Error> {
let name: TopicName = name.into();
let mut handles = self.create_topics(std::iter::once((name, opts)), backend)?;
Ok(handles
.pop()
.expect("single declaration must create one topic handle"))
}
pub fn create_topics(
&self,
declarations: impl IntoIterator<Item = (TopicName, TopicOptions)>,
backend: impl TopicBackend<T>,
) -> Result<Vec<TopicHandle<T>>, Error> {
let declarations: Vec<(TopicName, TopicOptions)> = declarations.into_iter().collect();
let mut topics = self.inner.topics.write();
let mut seen = HashSet::with_capacity(declarations.len());
for (name, _) in &declarations {
if topics.contains_key(name) || !seen.insert(name.clone()) {
return Err(Error::TopicAlreadyExists {
topic: name.clone(),
});
}
}
let mut handles = Vec::with_capacity(declarations.len());
for (name, opts) in declarations {
let state = backend.create_topic(name.clone(), opts);
_ = topics.insert(name, state.clone());
handles.push(TopicHandle::new(state));
}
Ok(handles)
}
pub fn create_in_memory_topic(
&self,
name: impl Into<TopicName>,
opts: TopicOptions,
) -> Result<TopicHandle<T>, Error> {
self.create_topic(name, opts, InMemoryBackend)
}
pub fn create_in_memory_topics(
&self,
declarations: impl IntoIterator<Item = (TopicName, TopicOptions)>,
) -> Result<Vec<TopicHandle<T>>, Error> {
self.create_topics(declarations, InMemoryBackend)
}
pub fn get_topic(&self, name: impl AsRef<str>) -> Option<TopicHandle<T>> {
let name = name.as_ref();
let topics = self.inner.topics.read();
topics
.get(name)
.map(|inner| TopicHandle::new(inner.clone()))
}
pub fn get_topic_required(&self, name: impl AsRef<str>) -> Result<TopicHandle<T>, Error> {
let name = name.as_ref();
self.get_topic(name).ok_or_else(|| Error::UnknownTopic {
topic: name.to_owned(),
})
}
pub fn has_topic(&self, name: impl AsRef<str>) -> bool {
let name = name.as_ref();
let topics = self.inner.topics.read();
topics.contains_key(name)
}
pub fn remove_topic(&self, name: impl AsRef<str>) -> bool {
let name = name.as_ref();
let mut topics = self.inner.topics.write();
if let Some(inner) = topics.remove(name) {
inner.close();
true
} else {
false
}
}
#[must_use]
pub fn topic_names(&self) -> Vec<TopicName> {
let topics = self.inner.topics.read();
topics.keys().cloned().collect()
}
pub fn close_all(&self) {
let mut topics = self.inner.topics.write();
for (_, inner) in topics.drain() {
inner.close();
}
}
}
impl<T: Send + Sync + 'static> Default for TopicBroker<T> {
fn default() -> Self {
Self::new()
}
}