use crate::Interests;
use crate::context_declaration::CompiledHeaderPropagationPolicy as HeaderPropagationPolicy;
use crate::control::{AckMsg, NackMsg};
use crate::effect_handler::{EffectHandlerCore, TelemetryTimerCancelHandle, TimerCancelHandle};
use crate::error::Error;
use crate::message::ExporterInbox;
use crate::node::NodeId;
use crate::runtime_services::{CodecEffectHandler, PipelineRuntimeServices};
use crate::terminal_state::TerminalState;
use async_trait::async_trait;
use otel_arrow_dfe_pdata_codec::CodecService;
use otel_arrow_dfe_telemetry::error::Error as TelemetryError;
use otel_arrow_dfe_telemetry::metrics::{MetricSet, MetricSetHandler};
use otel_arrow_dfe_telemetry::reporter::MetricsReporter;
use std::marker::PhantomData;
use std::rc::Rc;
use std::time::Duration;
#[async_trait( ? Send)]
pub trait Exporter<PData> {
async fn start(
self: Box<Self>,
inbox: ExporterInbox<PData>,
effect_handler: EffectHandler<PData>,
) -> Result<TerminalState, Error>;
}
#[derive(Clone)]
pub struct EffectHandler<PData> {
pub(crate) core: EffectHandlerCore<PData>,
_pd: PhantomData<PData>,
propagation_policy: Option<Rc<HeaderPropagationPolicy>>,
}
impl<PData> EffectHandler<PData> {
#[must_use]
pub fn new(
node_id: NodeId,
metrics_reporter: MetricsReporter,
runtime_services: PipelineRuntimeServices,
) -> Self {
EffectHandler {
core: EffectHandlerCore::new(node_id, metrics_reporter, runtime_services),
_pd: PhantomData,
propagation_policy: None,
}
}
#[must_use]
pub fn exporter_id(&self) -> NodeId {
self.core.node_id()
}
#[must_use]
pub fn node_interests(&self) -> Interests {
self.core.node_interests()
}
#[must_use]
pub fn propagation_policy(&self) -> Option<&HeaderPropagationPolicy> {
self.propagation_policy.as_deref()
}
pub fn set_propagation_policy(&mut self, policy: Option<HeaderPropagationPolicy>) {
self.propagation_policy = policy.map(Rc::new);
}
pub async fn info(&self, message: &str) {
self.core.info(message).await;
}
pub async fn start_periodic_timer(
&self,
duration: Duration,
) -> Result<TimerCancelHandle<PData>, Error> {
self.core.start_periodic_timer(duration).await
}
pub async fn start_periodic_telemetry(
&self,
duration: Duration,
) -> Result<TelemetryTimerCancelHandle<PData>, Error> {
self.core.start_periodic_telemetry(duration).await
}
#[allow(dead_code)] pub(crate) fn report_metrics<M: MetricSetHandler + 'static>(
&mut self,
metrics: &mut MetricSet<M>,
) -> Result<(), TelemetryError> {
self.core.report_metrics(metrics)
}
pub fn set_pipeline_completion_msg_sender(
&mut self,
pipeline_completion_msg_sender: crate::control::PipelineCompletionMsgSender<PData>,
) {
self.core
.set_pipeline_completion_msg_sender(pipeline_completion_msg_sender);
}
}
impl<PData> CodecEffectHandler for EffectHandler<PData> {
fn codec_service(&self) -> &CodecService {
self.core.runtime_services.codecs()
}
}
#[async_trait(?Send)]
impl<PData: crate::Unwindable> crate::_private::AckNackRouting<PData> for EffectHandler<PData> {
async fn route_ack(&self, ack: AckMsg<PData>) -> Result<(), Error> {
self.core.route_ack(ack).await
}
async fn route_nack(&self, nack: NackMsg<PData>) -> Result<(), Error> {
self.core.route_nack(nack).await
}
}