#![feature(is_some_and)]
#![warn(missing_docs)]
use serde::Serialize;
use std::{
env,
net::{SocketAddr, UdpSocket},
result::Result as StdResult,
sync::Arc,
};
mod epoch;
mod error;
mod header;
mod hexbytes;
mod lambda;
pub mod segment;
mod segment_id;
mod trace_id;
mod tracing;
pub use crate::{
epoch::Seconds,
error::Error,
header::Header,
segment::*,
segment_id::SegmentId,
trace_id::TraceId,
tracing::XRaySubscriber,
tracing::aws_metadata,
};
pub type Result<T> = StdResult<T, Error>;
#[derive(Debug)]
pub struct Client {
socket: Arc<UdpSocket>,
}
impl Default for Client {
fn default() -> Self {
let addr: SocketAddr = env::var("AWS_XRAY_DAEMON_ADDRESS")
.ok()
.and_then(|value| value.parse::<SocketAddr>().ok())
.unwrap_or_else(|| {
log::trace!("No valid `AWS_XRAY_DAEMON_ADDRESS` env variable detected falling back on default: 127.0.0.1:2000");
([127, 0, 0, 1], 2000).into()
});
Client::new(addr).expect("failed to connect to socket")
}
}
impl Client {
const HEADER: &'static [u8] = br#"{"format": "json", "version": 1}
"#;
pub fn new(addr: SocketAddr) -> Result<Self> {
let socket = Arc::new(UdpSocket::bind(&[([0, 0, 0, 0], 0).into()][..])?);
socket.set_nonblocking(true)?;
socket.connect(&addr)?;
log::trace!("connecting to xray daemon {}", addr);
Ok(Client { socket })
}
#[inline]
fn packet<S>(data: S) -> Result<Vec<u8>>
where
S: Serialize,
{
let bytes = serde_json::to_vec(&data)?;
Ok([Self::HEADER, &bytes].concat())
}
pub fn send<S>(
&self,
data: &S,
) -> Result<()>
where
S: Serialize,
{
log::trace!(
"sending trace data {}",
serde_json::to_string_pretty(&data).unwrap_or_default()
);
let out = self.socket.send(&Self::packet(data)?)?;
log::trace!("send? {:?}", out);
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
#[ignore]
fn client_can_send_data() {
env_logger::init();
let mut segment = Segment::begin(
"test-segment",
SegmentId::default(),
None,
TraceId::default(),
);
std::thread::sleep(std::time::Duration::from_secs(1));
segment.end();
if let Err(e) = Client::default().send(&segment) {
assert!(false, "failed to send data: {}", e)
}
}
#[test]
fn client_prefixes_packets_with_header() {
assert_eq!(
Client::packet(serde_json::json!({
"foo": "bar"
}))
.unwrap(),
br#"{"format": "json", "version": 1}
{"foo":"bar"}"#
.to_vec()
)
}
}