mod otlp;
mod reporter;
#[cfg(all(test, feature = "test-broker"))]
mod tests;
use std::sync::Arc;
use std::time::Duration;
use tokio::sync::watch;
use tokio::task::JoinHandle;
use crate::client::Kafka;
use crate::error::{KrafkaError, Result};
use crate::metrics::{ClientInstanceId, MetricsSource};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum ClientType {
Producer,
Consumer,
Admin,
}
impl ClientType {
fn name(self) -> &'static str {
match self {
Self::Producer => "producer",
Self::Consumer => "consumer",
Self::Admin => "admin",
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum InstanceId {
Pending,
Assigned(ClientInstanceId),
Unavailable,
}
#[derive(Debug)]
pub(crate) struct Telemetry(Option<Running>);
#[derive(Debug)]
struct Running {
stop: watch::Sender<bool>,
instance_id: watch::Receiver<InstanceId>,
task: parking_lot::Mutex<Option<JoinHandle<()>>>,
}
impl Telemetry {
pub(crate) fn disabled() -> Self {
Self(None)
}
pub(crate) fn start(
enabled: bool,
kafka: &Kafka,
client_type: ClientType,
source: Arc<MetricsSource>,
) -> Self {
if !enabled {
return Self::disabled();
}
let Ok(runtime) = tokio::runtime::Handle::try_current() else {
tracing::debug!("no Tokio runtime; client telemetry not started");
return Self::disabled();
};
let (stop, stop_rx) = watch::channel(false);
let (id_tx, instance_id) = watch::channel(InstanceId::Pending);
let reporter = reporter::Reporter::new(kafka.clone(), client_type, source);
let task = runtime.spawn(reporter.run(stop_rx, id_tx));
Self(Some(Running {
stop,
instance_id,
task: parking_lot::Mutex::new(Some(task)),
}))
}
pub(crate) async fn close(&self, timeout: Duration) {
let Some(running) = &self.0 else {
return;
};
let Some(task) = running.task.lock().take() else {
return;
};
let _ = running.stop.send(true);
let mut task = AbortOnDrop(task);
if tokio::time::timeout(timeout, &mut task.0).await.is_err() {
tracing::debug!("client telemetry did not finish its terminating push in time");
}
}
pub(crate) async fn client_instance_id(
&self,
timeout: Duration,
) -> Result<Option<ClientInstanceId>> {
let Some(running) = &self.0 else {
return Err(KrafkaError::illegal_state(
"client telemetry is disabled; enable it with metrics_push(true)",
));
};
let mut rx = running.instance_id.clone();
let wait = rx.wait_for(|id| *id != InstanceId::Pending);
match tokio::time::timeout(timeout, wait).await {
Ok(Ok(id)) => Ok(match *id {
InstanceId::Assigned(id) => Some(id),
InstanceId::Pending | InstanceId::Unavailable => None,
}),
Ok(Err(_)) => Ok(None),
Err(_) => Err(KrafkaError::timeout("waiting for the client instance id")),
}
}
}
struct AbortOnDrop(JoinHandle<()>);
impl Drop for AbortOnDrop {
fn drop(&mut self) {
self.0.abort();
}
}
impl Drop for Telemetry {
fn drop(&mut self) {
if let Some(running) = &self.0
&& let Some(task) = running.task.lock().take()
{
task.abort();
}
}
}