use std::ops::Deref;
use crate::topic::handle::TopicHandle;
use otel_arrow_dfe_config::topic::{TopicAckPropagationMode, TopicQueueOnFullPolicy};
pub struct PipelineTopicBinding<T: Send + Sync + 'static> {
handle: TopicHandle<T>,
queue_on_full_default: TopicQueueOnFullPolicy,
ack_propagation_mode_default: TopicAckPropagationMode,
}
impl<T: Send + Sync + 'static> Clone for PipelineTopicBinding<T> {
fn clone(&self) -> Self {
Self {
handle: self.handle.clone(),
queue_on_full_default: self.queue_on_full_default.clone(),
ack_propagation_mode_default: self.ack_propagation_mode_default,
}
}
}
impl<T: Send + Sync + 'static> PipelineTopicBinding<T> {
#[must_use]
pub fn new(handle: TopicHandle<T>) -> Self {
Self {
handle,
queue_on_full_default: TopicQueueOnFullPolicy::Block,
ack_propagation_mode_default: TopicAckPropagationMode::Disabled,
}
}
#[must_use]
pub fn with_default_queue_on_full(&self, policy: TopicQueueOnFullPolicy) -> Self {
Self {
handle: self.handle.clone(),
queue_on_full_default: policy,
ack_propagation_mode_default: self.ack_propagation_mode_default,
}
}
#[must_use]
pub fn with_default_ack_propagation_mode(&self, mode: TopicAckPropagationMode) -> Self {
Self {
handle: self.handle.clone(),
queue_on_full_default: self.queue_on_full_default.clone(),
ack_propagation_mode_default: mode,
}
}
#[must_use]
pub const fn handle(&self) -> &TopicHandle<T> {
&self.handle
}
#[must_use]
pub fn into_handle(self) -> TopicHandle<T> {
self.handle
}
#[must_use]
pub fn default_queue_on_full(&self) -> TopicQueueOnFullPolicy {
self.queue_on_full_default.clone()
}
#[must_use]
pub const fn default_ack_propagation_mode(&self) -> TopicAckPropagationMode {
self.ack_propagation_mode_default
}
}
impl<T: Send + Sync + 'static> From<TopicHandle<T>> for PipelineTopicBinding<T> {
fn from(value: TopicHandle<T>) -> Self {
Self::new(value)
}
}
impl<T: Send + Sync + 'static> Deref for PipelineTopicBinding<T> {
type Target = TopicHandle<T>;
fn deref(&self) -> &Self::Target {
&self.handle
}
}