use std::collections::HashMap;
use async_trait::async_trait;
use crate::envelope::EventEnvelope;
use crate::function::{AppError, ComposableFunction};
use crate::platform::Platform;
pub const DISTRIBUTED_TRACING: &str = "distributed.tracing";
pub const DISTRIBUTED_TRACE_FORWARDER: &str = "distributed.trace.forwarder";
pub const TRANSACTION_JOURNAL_RECORDER: &str = "transaction.journal.recorder";
pub const ZERO_TRACING_FILTER: [&str; 4] = [
DISTRIBUTED_TRACING,
DISTRIBUTED_TRACE_FORWARDER,
TRANSACTION_JOURNAL_RECORDER,
crate::inbox::TEMPORARY_INBOX,
];
pub struct Telemetry {
platform: Platform,
}
impl Telemetry {
pub fn new(platform: &Platform) -> Self {
Telemetry {
platform: platform.clone(),
}
}
}
#[async_trait]
impl ComposableFunction for Telemetry {
async fn handle_event(
&self,
_headers: HashMap<String, String>,
input: EventEnvelope,
_instance: usize,
) -> Result<EventEnvelope, AppError> {
let Ok(payload) = input.body_as::<serde_json::Value>() else {
return Ok(EventEnvelope::new());
};
let Some(payload) = payload.as_object() else {
return Ok(EventEnvelope::new());
};
let mut metrics = match payload.get("trace").and_then(|t| t.as_object()) {
Some(m) if !m.is_empty() => m.clone(),
_ => return Ok(EventEnvelope::new()),
};
let Some(service) = permitted_route(metrics.get("service")) else {
return Ok(EventEnvelope::new());
};
metrics.insert("service".to_string(), serde_json::Value::String(service));
if let Some(from) = metrics.get("from").and_then(|f| f.as_str()) {
let trimmed = trim_origin(from).to_string();
metrics.insert("from".to_string(), serde_json::Value::String(trimmed));
}
let annotations = payload
.get("annotations")
.and_then(|a| a.as_object())
.cloned()
.unwrap_or_default();
let mut dataset = serde_json::Map::new();
dataset.insert("trace".to_string(), serde_json::Value::Object(metrics));
if !annotations.is_empty() {
dataset.insert(
"annotations".to_string(),
serde_json::Value::Object(annotations),
);
}
let dataset = serde_json::Value::Object(dataset);
log::info!("{dataset}");
if self.platform.has_route(DISTRIBUTED_TRACE_FORWARDER) {
let event = EventEnvelope::new()
.set_to(DISTRIBUTED_TRACE_FORWARDER)
.set_body(&dataset)?;
let _ = self
.platform
.deliver(DISTRIBUTED_TRACE_FORWARDER, event)
.await;
}
if payload.contains_key("journal") && self.platform.has_route(TRANSACTION_JOURNAL_RECORDER)
{
let mut forward = dataset;
if let (Some(map), Some(journal)) = (forward.as_object_mut(), payload.get("journal")) {
map.insert("journal".to_string(), journal.clone());
}
let event = EventEnvelope::new()
.set_to(TRANSACTION_JOURNAL_RECORDER)
.set_body(&forward)?;
let _ = self
.platform
.deliver(TRANSACTION_JOURNAL_RECORDER, event)
.await;
}
Ok(EventEnvelope::new())
}
}
fn permitted_route(service: Option<&serde_json::Value>) -> Option<String> {
let route = service?.as_str()?;
let name = trim_origin(route);
if ZERO_TRACING_FILTER.contains(&name) {
None
} else {
Some(name.to_string())
}
}
fn trim_origin(route: &str) -> &str {
match route.find('@') {
Some(at) => &route[..at],
None => route,
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn plumbing_routes_are_filtered() {
let v = |s: &str| serde_json::Value::String(s.to_string());
assert_eq!(
permitted_route(Some(&v("v1.hello"))),
Some("v1.hello".into())
);
assert_eq!(
permitted_route(Some(&v("v1.hello@abc123"))),
Some("v1.hello".into())
);
assert_eq!(permitted_route(Some(&v("distributed.tracing"))), None);
assert_eq!(
permitted_route(Some(&v("distributed.trace.forwarder@x"))),
None
);
assert_eq!(permitted_route(None), None);
}
}