use rmp::{ decode::{Bytes, MarkerReadError, RmpRead, bytes::BytesReadError}, encode::{ByteBuf, RmpWrite}, }; use tracing::{Event, field::Visit, span}; use crate::logs::{AttributeValue, Attributes}; #[derive(Debug, Clone)] pub struct MessagePackBytes(Vec); impl MessagePackBytes { fn bytes(&self) -> Bytes<'_> { Bytes::new(&self.0[..]) } } impl From<&span::Attributes<'_>> for MessagePackBytes { fn from(value: &span::Attributes<'_>) -> Self { let metadata = value.metadata(); let mut extra_attrs = 3; // name, target, level if metadata.module_path().is_some() { extra_attrs += 1; } if metadata.file().is_some() { extra_attrs += 1; } if metadata.line().is_some() { extra_attrs += 1; } let mut visitor = Visitor::new(value.fields().len() + extra_attrs); // TODO: maybe allow controlling which of these actually get recorded somehow? visitor.raw_record_str("name", metadata.name()); visitor.raw_record_str("target", metadata.target()); visitor.raw_record_str("level", metadata.level().as_str()); // TODO: as int? if let Some(module_path) = metadata.module_path() { visitor.raw_record_str("module_path", module_path); } if let Some(file) = metadata.file() { visitor.raw_record_str("file", file); } if let Some(line) = metadata.line() { visitor.raw_record_u64("line", u64::from(line)); } value.record(&mut visitor); visitor.into() } } impl From<&span::Record<'_>> for MessagePackBytes { fn from(value: &span::Record<'_>) -> Self { let mut visitor = Visitor::new(value.len()); value.record(&mut visitor); visitor.into() } } impl From<&Event<'_>> for MessagePackBytes { fn from(value: &Event<'_>) -> Self { let mut visitor = Visitor::new(value.fields().count()); value.record(&mut visitor); visitor.into() } } const SIGNED_128: i8 = 27; const UNSIGNED_128: i8 = 34; #[derive(Debug, Default)] pub(crate) struct Visitor(ByteBuf); impl Visitor { pub(crate) fn new(len: usize) -> Self { let mut this = Self(ByteBuf::new()); rmp::encode::write_map_len(&mut this.0, len.try_into().unwrap()).unwrap(); this } } // the actual bulk of the record bits. these are pub(crate) so we can use them when testing - // tracing spans are a pain / to construct impl Visitor { pub(crate) fn raw_record_field_name_str(&mut self, field: &'static str) { rmp::encode::write_str(&mut self.0, field).unwrap(); } pub(crate) fn raw_record_debug(&mut self, field: &'static str, value: &dyn core::fmt::Debug) { self.raw_record_field_name_str(field); let s = format!("{value:?}"); rmp::encode::write_str(&mut self.0, &s).unwrap(); } pub(crate) fn raw_record_f64(&mut self, field: &'static str, value: f64) { self.raw_record_field_name_str(field); rmp::encode::write_f64(&mut self.0, value).unwrap(); } pub(crate) fn raw_record_i64(&mut self, field: &'static str, value: i64) { self.raw_record_field_name_str(field); rmp::encode::write_i64(&mut self.0, value).unwrap(); } pub(crate) fn raw_record_u64(&mut self, field: &'static str, value: u64) { self.raw_record_field_name_str(field); rmp::encode::write_u64(&mut self.0, value).unwrap(); } pub(crate) fn raw_record_i128(&mut self, field: &'static str, value: i128) { self.raw_record_field_name_str(field); rmp::encode::write_ext_meta(&mut self.0, 16, SIGNED_128).unwrap(); self.0.write_bytes(&value.to_be_bytes()).unwrap(); } pub(crate) fn raw_record_u128(&mut self, field: &'static str, value: u128) { self.raw_record_field_name_str(field); rmp::encode::write_ext_meta(&mut self.0, 16, UNSIGNED_128).unwrap(); self.0.write_bytes(&value.to_be_bytes()).unwrap(); } pub(crate) fn raw_record_bool(&mut self, field: &'static str, value: bool) { self.raw_record_field_name_str(field); rmp::encode::write_bool(&mut self.0, value).unwrap(); } pub(crate) fn raw_record_str(&mut self, field: &'static str, value: &str) { self.raw_record_field_name_str(field); rmp::encode::write_str(&mut self.0, value).unwrap(); } pub(crate) fn raw_record_bytes(&mut self, field: &'static str, value: &[u8]) { self.raw_record_field_name_str(field); rmp::encode::write_bin(&mut self.0, value).unwrap(); } } impl Visit for Visitor { fn record_debug(&mut self, field: &tracing::field::Field, value: &dyn core::fmt::Debug) { self.raw_record_debug(field.name(), value); } fn record_f64(&mut self, field: &tracing::field::Field, value: f64) { self.raw_record_f64(field.name(), value); } fn record_i64(&mut self, field: &tracing::field::Field, value: i64) { self.raw_record_i64(field.name(), value); } fn record_u64(&mut self, field: &tracing::field::Field, value: u64) { self.raw_record_u64(field.name(), value); } fn record_i128(&mut self, field: &tracing::field::Field, value: i128) { self.raw_record_i128(field.name(), value); } fn record_u128(&mut self, field: &tracing::field::Field, value: u128) { self.raw_record_u128(field.name(), value); } fn record_bool(&mut self, field: &tracing::field::Field, value: bool) { self.raw_record_bool(field.name(), value); } fn record_str(&mut self, field: &tracing::field::Field, value: &str) { self.raw_record_str(field.name(), value); } fn record_bytes(&mut self, field: &tracing::field::Field, value: &[u8]) { self.raw_record_bytes(field.name(), value); } } impl From for MessagePackBytes { fn from(vis: Visitor) -> Self { Self(vis.0.into_vec()) } } impl Visitor {} pub(crate) const INTERNAL_MESSAGEPACK_ERROR: &str = "invalid messagepack data received internally"; impl TryInto for &MessagePackBytes { type Error = rmp::decode::ValueReadError; fn try_into(self) -> Result { let mut rd = self.bytes(); let mut attrs = Vec::new(); // note: the .clone means that we use a temporary cursor, ie we peek rather than consuming // this is needed because the read_ methods expect to consume the marker, but we need // to know what it is first (and, in some cases, allocate a buffer of the correct size). let len = rmp::decode::read_map_len(&mut rd)?; for _ in 0..len { let name_len = rmp::decode::read_str_len(&mut rd.clone())?; let mut name = vec![0; name_len.try_into().unwrap()]; rmp::decode::read_str(&mut rd, &mut name).expect(INTERNAL_MESSAGEPACK_ERROR); let mut after_marker = rd; let value = match rmp::decode::read_marker(&mut after_marker)? { rmp::Marker::F64 => AttributeValue::F64(rmp::decode::read_f64(&mut rd)?), rmp::Marker::U64 => AttributeValue::U64(rmp::decode::read_u64(&mut rd)?), rmp::Marker::I64 => AttributeValue::I64(rmp::decode::read_i64(&mut rd)?), rmp::Marker::False => AttributeValue::Bool(false), rmp::Marker::True => AttributeValue::Bool(true), rmp::Marker::FixStr(_) | rmp::Marker::Str8 | rmp::Marker::Str16 | rmp::Marker::Str32 => { let len = rmp::decode::read_str_len(&mut rd.clone())? as usize; let mut buf = vec![0; len]; rmp::decode::read_str(&mut rd, &mut buf).expect(INTERNAL_MESSAGEPACK_ERROR); AttributeValue::String( String::from_utf8(buf).expect(INTERNAL_MESSAGEPACK_ERROR), ) } rmp::Marker::FixExt16 => { let signed = match rmp::decode::read_i8(&mut after_marker)? { SIGNED_128 => true, UNSIGNED_128 => false, _ => unreachable!(), }; let mut buf = [0; 16]; rd.read_exact_buf(&mut buf).map_err(MarkerReadError)?; if signed { AttributeValue::U128(u128::from_be_bytes(buf)) } else { AttributeValue::I128(i128::from_be_bytes(buf)) } } rmp::Marker::Bin8 | rmp::Marker::Bin16 | rmp::Marker::Bin32 => { let len = rmp::decode::read_bin_len(&mut rd)? as usize; let mut buf = vec![0; len]; rd.read_exact_buf(&mut buf).map_err(MarkerReadError)?; AttributeValue::Bytes(buf.into_boxed_slice()) } rmp::Marker::FixPos(_) | rmp::Marker::FixMap(_) | rmp::Marker::FixArray(_) | rmp::Marker::Null | rmp::Marker::Reserved | rmp::Marker::Ext8 | rmp::Marker::Ext16 | rmp::Marker::Ext32 | rmp::Marker::F32 | rmp::Marker::U8 | rmp::Marker::U16 | rmp::Marker::U32 | rmp::Marker::I8 | rmp::Marker::I16 | rmp::Marker::I32 | rmp::Marker::FixExt1 | rmp::Marker::FixExt2 | rmp::Marker::FixExt4 | rmp::Marker::FixExt8 | rmp::Marker::Array16 | rmp::Marker::Array32 | rmp::Marker::Map16 | rmp::Marker::Map32 | rmp::Marker::FixNeg(_) => unreachable!(), }; attrs.push(( String::from_utf8(name).expect(INTERNAL_MESSAGEPACK_ERROR), value, )); } Ok(attrs) } }