use std::borrow::Cow;
use tracing::Span;
use crate::error::KrafkaError;
pub const OTEL_SEMCONV_VERSION: &str = "1.44.0";
#[inline]
pub(crate) fn send_span(topic: &str, client_id: &str, key: Option<&[u8]>, tombstone: bool) -> Span {
tracing::info_span!(
target: "krafka::producer",
"send",
otel.name = %format_args!("send {topic}"),
otel.kind = "producer",
otel.status_code = tracing::field::Empty,
messaging.system = "kafka",
messaging.operation.name = "send",
messaging.operation.type = "send",
messaging.destination.name = topic,
messaging.destination.partition.id = tracing::field::Empty,
messaging.kafka.offset = tracing::field::Empty,
messaging.kafka.message.tombstone = tombstone.then_some(true),
messaging.client.id = client_id,
krafka.message.key.size = key.map(|k| k.len() as u64),
error.type = tracing::field::Empty,
)
}
#[inline]
pub(crate) fn poll_span(group: Option<&str>, client_id: &str) -> Span {
tracing::info_span!(
target: "krafka::consumer",
"poll",
otel.name = "poll",
otel.kind = "client",
otel.status_code = tracing::field::Empty,
messaging.system = "kafka",
messaging.operation.name = "poll",
messaging.operation.type = "receive",
messaging.destination.name = tracing::field::Empty,
messaging.consumer.group.name = group,
messaging.client.id = client_id,
messaging.batch.message_count = tracing::field::Empty,
error.type = tracing::field::Empty,
)
}
#[inline]
pub(crate) fn commit_span(group: Option<&str>, client_id: &str) -> Span {
tracing::info_span!(
target: "krafka::consumer",
"commit",
otel.name = "commit",
otel.kind = "client",
otel.status_code = tracing::field::Empty,
messaging.system = "kafka",
messaging.operation.name = "commit",
messaging.operation.type = "settle",
messaging.destination.name = tracing::field::Empty,
messaging.consumer.group.name = group,
messaging.client.id = client_id,
error.type = tracing::field::Empty,
)
}
#[inline]
pub(crate) fn rebalance_span(group: &str, protocol: &str, generation: Option<i32>) -> Span {
tracing::info_span!(
target: "krafka::consumer",
"rebalance",
otel.name = %format_args!("rebalance {group}"),
otel.kind = "internal",
messaging.system = "kafka",
messaging.consumer.group.name = group,
krafka.rebalance.protocol = protocol,
krafka.rebalance.generation = generation,
krafka.rebalance.assigned = tracing::field::Empty,
krafka.rebalance.revoked = tracing::field::Empty,
krafka.rebalance.partitions = tracing::field::Empty,
)
}
pub(crate) fn record_destination<'a>(
span: &Span,
operation: &'static str,
topics: impl IntoIterator<Item = &'a str>,
) {
let mut topics = topics.into_iter();
if let (Some(topic), None) = (topics.next(), topics.next()) {
span.record(
"otel.name",
tracing::field::display(OperationName(operation, Some(topic))),
);
span.record("messaging.destination.name", topic);
}
}
struct OperationName<'a>(&'static str, Option<&'a str>);
impl std::fmt::Display for OperationName<'_> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self.1 {
Some(topic) => write!(f, "{} {topic}", self.0),
None => f.write_str(self.0),
}
}
}
pub(crate) fn record_poll_outcome(span: &Span, outcome: std::result::Result<usize, &KrafkaError>) {
match outcome {
Ok(count) => {
span.record("messaging.batch.message_count", count as u64);
}
Err(error) => record_error(span, error),
}
}
pub(crate) fn record_error(span: &Span, error: &KrafkaError) {
if !span.is_disabled() {
span.record("error.type", error_type(error).as_ref());
span.record("otel.status_code", "error");
}
}
pub(crate) fn error_type(error: &KrafkaError) -> Cow<'static, str> {
match error {
KrafkaError::Broker { code, .. } => Cow::Owned(screaming_snake(&format!("{code:?}"))),
KrafkaError::OutOfOrderSequence { .. } => "OUT_OF_ORDER_SEQUENCE_NUMBER".into(),
KrafkaError::UnknownTopic { .. } => "UNKNOWN_TOPIC_OR_PARTITION".into(),
KrafkaError::Network(_) => "network".into(),
KrafkaError::Protocol { .. } => "protocol".into(),
KrafkaError::Auth { .. } => "auth".into(),
KrafkaError::Timeout { .. } => "timeout".into(),
KrafkaError::DeliveryTimeout { .. } => "delivery_timeout".into(),
KrafkaError::Config { .. } => "config".into(),
KrafkaError::Compression { .. } => "compression".into(),
KrafkaError::Serialization { .. } => "serialization".into(),
KrafkaError::RecordDeserialization { .. } => "record_deserialization".into(),
KrafkaError::Closed { .. } => "closed".into(),
KrafkaError::Wakeup => "wakeup".into(),
KrafkaError::Fenced { .. } => "fenced".into(),
KrafkaError::TransactionAbortable { .. } => "transaction_abortable".into(),
KrafkaError::NoOffset { .. } => "no_offset".into(),
KrafkaError::IllegalState { .. } => "illegal_state".into(),
}
}
fn screaming_snake(camel: &str) -> String {
let name = camel.split('(').next().unwrap_or(camel);
let mut out = String::with_capacity(name.len() + 8);
for (i, ch) in name.chars().enumerate() {
if ch.is_ascii_uppercase() && i > 0 {
out.push('_');
}
out.push(ch.to_ascii_uppercase());
}
out
}
#[cfg(all(test, feature = "test-broker"))]
mod broker_tests;
#[cfg(test)]
#[allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)]
mod tests {
use super::*;
use crate::error::ErrorCode;
#[test]
fn broker_errors_are_typed_by_their_kafka_name() {
let error = KrafkaError::broker(ErrorCode::NotEnoughReplicas, "x");
assert_eq!(error_type(&error), "NOT_ENOUGH_REPLICAS");
let error = KrafkaError::broker(ErrorCode::UnknownTopicOrPartition, "x");
assert_eq!(error_type(&error), "UNKNOWN_TOPIC_OR_PARTITION");
assert_eq!(error_type(&KrafkaError::closed("x")), "closed");
}
#[test]
fn the_operation_name_omits_a_missing_destination() {
assert_eq!(
OperationName("poll", Some("orders")).to_string(),
"poll orders"
);
assert_eq!(OperationName("poll", None).to_string(), "poll");
}
}