use std::fmt;
use std::sync::{Arc, OnceLock};
use bytes::Bytes;
use ruststream::testing::{Coordinator, TestableBroker};
use ruststream::{Broker, DescribeServer, OutgoingMessage, RawMessage, ServerSpec, Subscribe};
use super::publisher::KafkaTestPublisher;
use super::router::KeyRouter;
use super::subscriber::KafkaTestSubscriber;
use crate::error::KafkaError;
pub(crate) struct TestBrokerState {
pub(crate) router: KeyRouter,
coordinator: OnceLock<Coordinator>,
}
impl TestBrokerState {
pub(crate) fn install(&self, coordinator: Coordinator) {
let _ = self.coordinator.set(coordinator);
}
pub(crate) fn coordinator(&self) -> Option<Coordinator> {
self.coordinator.get().cloned()
}
}
impl Default for TestBrokerState {
fn default() -> Self {
Self {
router: KeyRouter::default(),
coordinator: OnceLock::new(),
}
}
}
impl fmt::Debug for TestBrokerState {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("TestBrokerState")
.field("router", &self.router)
.finish_non_exhaustive()
}
}
#[derive(Debug, Clone, Default)]
pub struct KafkaTestBroker {
state: Arc<TestBrokerState>,
}
impl KafkaTestBroker {
#[must_use]
pub fn new() -> Self {
Self::default()
}
#[allow(clippy::unused_async)]
pub async fn subscribe(
&self,
topic: impl Into<String>,
) -> Result<KafkaTestSubscriber, KafkaError> {
self.subscribe_topics(std::slice::from_ref(&topic.into()))
.await
}
#[allow(clippy::unused_async)]
pub async fn subscribe_topics(
&self,
topics: &[String],
) -> Result<KafkaTestSubscriber, KafkaError> {
for topic in topics {
if topic.is_empty() {
return Err(KafkaError::InvalidOptions(
"topic name must not be empty; subscribe with the topic the handler \
consumes"
.to_owned(),
));
}
if topic.starts_with('^') {
return Err(KafkaError::InvalidOptions(format!(
"the in-process test broker routes by exact topic name; the pattern \
{topic:?} needs a real cluster",
)));
}
}
Ok(KafkaTestSubscriber::open_many(&self.state, topics))
}
#[must_use]
pub fn publisher(&self) -> KafkaTestPublisher {
KafkaTestPublisher::new(Arc::clone(&self.state))
}
}
impl Broker for KafkaTestBroker {
type Error = KafkaError;
async fn connect(&self) -> Result<(), Self::Error> {
Ok(())
}
async fn shutdown(&self) -> Result<(), Self::Error> {
self.state.router.clear();
Ok(())
}
}
#[allow(clippy::use_self)]
impl Subscribe for KafkaTestBroker {
type Subscriber = KafkaTestSubscriber;
async fn subscribe(&self, name: &str) -> Result<Self::Subscriber, Self::Error> {
KafkaTestBroker::subscribe(self, name).await
}
}
impl DescribeServer for KafkaTestBroker {
fn describe_server(&self) -> ServerSpec {
ServerSpec::in_process("kafka")
}
}
impl TestableBroker for KafkaTestBroker {
fn install_coordinator(&self, coordinator: Coordinator) {
self.state.install(coordinator);
}
fn inject(&self, message: OutgoingMessage<'_>) {
self.state.router.publish(
message.name(),
&Bytes::copy_from_slice(message.payload()),
message.headers(),
self.state.coordinator().as_ref(),
);
}
fn published(&self, name: &str) -> Vec<RawMessage> {
self.state.router.published(name)
}
}
ruststream::register_testable_broker!(KafkaTestBroker);