use rdkafka::TopicPartitionList;
use rdkafka::client::ClientContext;
use rdkafka::consumer::{BaseConsumer, Consumer, ConsumerContext};
use rdkafka::statistics::Statistics;
use rdkafka::types::RDKafkaRespErr;
use std::collections::VecDeque;
use std::sync::Mutex;
use std::sync::atomic::{AtomicBool, Ordering};
#[derive(Debug)]
pub(crate) enum Intent {
Assign(TopicPartitionList),
Revoke(TopicPartitionList),
Error(String),
}
#[derive(Debug, Default)]
pub(crate) struct SourceContext {
pub(crate) intents: Mutex<VecDeque<Intent>>,
pub(crate) stats: Mutex<Option<Box<Statistics>>>,
pub(crate) closing: AtomicBool,
}
impl ClientContext for SourceContext {
fn stats(&self, statistics: Statistics) {
*self.stats.lock().expect("stats lock") = Some(Box::new(statistics));
}
fn log(&self, level: rdkafka::config::RDKafkaLogLevel, fac: &str, log_message: &str) {
use rdkafka::config::RDKafkaLogLevel as L;
match level {
L::Emerg | L::Alert | L::Critical | L::Error => {
tracing::error!(target: "librdkafka", fac, "{log_message}");
}
L::Warning => tracing::warn!(target: "librdkafka", fac, "{log_message}"),
L::Notice | L::Info => tracing::info!(target: "librdkafka", fac, "{log_message}"),
L::Debug => tracing::debug!(target: "librdkafka", fac, "{log_message}"),
}
}
fn error(&self, error: rdkafka::error::KafkaError, reason: &str) {
tracing::warn!(target: "librdkafka", %error, "{reason}");
}
}
impl ConsumerContext for SourceContext {
fn rebalance(
&self,
base_consumer: &BaseConsumer<Self>,
err: RDKafkaRespErr,
tpl: &mut TopicPartitionList,
) {
if self.closing.load(Ordering::Acquire) {
let result = match err {
RDKafkaRespErr::RD_KAFKA_RESP_ERR__ASSIGN_PARTITIONS => base_consumer.assign(tpl),
_ => base_consumer.unassign(),
};
if let Err(e) = result {
tracing::warn!(error = %e, "rebalance completion during close failed");
}
return;
}
let intent = match err {
RDKafkaRespErr::RD_KAFKA_RESP_ERR__ASSIGN_PARTITIONS => Intent::Assign(tpl.clone()),
RDKafkaRespErr::RD_KAFKA_RESP_ERR__REVOKE_PARTITIONS => Intent::Revoke(tpl.clone()),
other => {
let code: rdkafka::error::RDKafkaErrorCode = other.into();
Intent::Error(code.to_string())
}
};
self.intents.lock().expect("intent lock").push_back(intent);
}
}