use std::fmt;
use std::sync::Arc;
use crate::consumer::{FetchedRecord, OffsetAndMetadata, TopicPartition};
use crate::error::Error;
use crate::producer::{ProduceRecord, RecordMetadata};
pub trait ProducerInterceptor: Send + Sync + 'static {
fn on_send(&self, rec: ProduceRecord) -> ProduceRecord {
rec
}
fn on_ack(&self, _md: &RecordMetadata) {}
fn on_error(&self, _err: &Error) {}
fn close(&self) {}
}
pub trait ConsumerInterceptor: Send + Sync + 'static {
fn on_consume(&self, recs: Vec<FetchedRecord>) -> Vec<FetchedRecord> {
recs
}
fn on_commit(&self, _offsets: &[(TopicPartition, OffsetAndMetadata)]) {}
fn close(&self) {}
}
#[derive(Clone, Default)]
pub struct ProducerInterceptors {
inner: Vec<Arc<dyn ProducerInterceptor>>,
}
impl ProducerInterceptors {
pub fn push(&mut self, i: impl ProducerInterceptor) {
self.inner.push(Arc::new(i));
}
pub(crate) fn on_send(&self, rec: ProduceRecord) -> ProduceRecord {
let mut rec = rec;
for i in &self.inner {
rec = i.on_send(rec);
}
rec
}
pub(crate) fn on_ack(&self, md: &RecordMetadata) {
for i in &self.inner {
i.on_ack(md);
}
}
pub(crate) fn on_error(&self, err: &Error) {
for i in &self.inner {
i.on_error(err);
}
}
pub(crate) fn close(&self) {
for i in &self.inner {
i.close();
}
}
}
impl fmt::Debug for ProducerInterceptors {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_tuple("ProducerInterceptors")
.field(&self.inner.len())
.finish()
}
}
#[derive(Clone, Default)]
pub struct ConsumerInterceptors {
inner: Vec<Arc<dyn ConsumerInterceptor>>,
}
impl ConsumerInterceptors {
pub fn push(&mut self, i: impl ConsumerInterceptor) {
self.inner.push(Arc::new(i));
}
pub(crate) fn on_consume(&self, recs: Vec<FetchedRecord>) -> Vec<FetchedRecord> {
let mut recs = recs;
for i in &self.inner {
recs = i.on_consume(recs);
}
recs
}
pub(crate) fn on_commit(&self, offsets: &[(TopicPartition, OffsetAndMetadata)]) {
for i in &self.inner {
i.on_commit(offsets);
}
}
pub(crate) fn close(&self) {
for i in &self.inner {
i.close();
}
}
}
impl fmt::Debug for ConsumerInterceptors {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_tuple("ConsumerInterceptors")
.field(&self.inner.len())
.finish()
}
}