use std::path::Path;
use std::time::Duration;
use serde::{Deserialize, Serialize};
use serde_json::{Map, Value};
use super::client::{QueuedEvent, TelemetryHandle};
pub const PATH: &str = "/admin/telemetry/events";
pub const MAX_EVENTS: usize = 50;
pub const MAX_BYTES: usize = 64 * 1024;
const MAX_NAME_LEN: usize = 64;
const DRAIN_WAIT: Duration = Duration::from_millis(25);
const CONNECT_TIMEOUT: Duration = Duration::from_millis(50);
const TOTAL_TIMEOUT: Duration = Duration::from_millis(100);
#[derive(Debug, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct HandoffBatch {
pub events: Vec<HandoffEvent>,
}
#[derive(Debug, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct HandoffEvent {
pub event: String,
pub properties: Map<String, Value>,
}
impl HandoffBatch {
pub fn validate(&self) -> Result<(), &'static str> {
if self.events.is_empty() || self.events.len() > MAX_EVENTS {
return Err("a hand-off carries between 1 and MAX_EVENTS events");
}
for event in &self.events {
if !is_event_name(&event.event) {
return Err("an event name is not a telemetry event name");
}
for key in ["distinct_id", "agent_id"] {
if !event
.properties
.get(key)
.and_then(Value::as_str)
.is_some_and(|v| !v.is_empty())
{
return Err("an event is missing its distinct_id or agent_id");
}
}
}
Ok(())
}
pub fn into_queued(self) -> impl Iterator<Item = QueuedEvent> {
self.events.into_iter().map(|e| QueuedEvent {
name: e.event,
properties: e.properties,
})
}
}
fn is_event_name(name: &str) -> bool {
let bare = name.strip_prefix('$').unwrap_or(name);
name.len() <= MAX_NAME_LEN
&& bare.starts_with(|c: char| c.is_ascii_lowercase())
&& bare
.chars()
.all(|c| c.is_ascii_lowercase() || c.is_ascii_digit() || c == '_')
}
pub fn hand_off(handle: &TelemetryHandle, openlatch_dir: &Path) {
if !openlatch_dir.join("daemon.pid").exists() {
return;
}
let Ok(cfg) = crate::config::Config::load(None, None, false) else {
return;
};
let events = handle.take_unsent(DRAIN_WAIT);
if !events.is_empty() {
send(events, openlatch_dir, cfg.port);
}
}
fn send(events: Vec<QueuedEvent>, openlatch_dir: &Path, port: u16) -> bool {
let Ok(token) = std::fs::read_to_string(openlatch_dir.join("daemon.token")) else {
return false;
};
let token = token.trim();
if token.is_empty() {
return false;
}
let batch = HandoffBatch {
events: events
.into_iter()
.take(MAX_EVENTS)
.map(|e| HandoffEvent {
event: e.name,
properties: e.properties,
})
.collect(),
};
let Ok(body) = serde_json::to_vec(&batch) else {
return false;
};
if body.len() > MAX_BYTES {
return false;
}
let agent = ureq::Agent::new_with_config(
ureq::Agent::config_builder()
.timeout_connect(Some(CONNECT_TIMEOUT))
.timeout_global(Some(TOTAL_TIMEOUT))
.http_status_as_error(false)
.proxy(None)
.max_redirects(0)
.build(),
);
agent
.post(format!("http://127.0.0.1:{port}{PATH}"))
.header("Authorization", format!("Bearer {token}"))
.header("Content-Type", "application/json")
.send(&body[..])
.is_ok_and(|resp| resp.status().is_success())
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
fn event(name: &str) -> HandoffEvent {
let mut properties = Map::new();
properties.insert("distinct_id".into(), json!("agt_a"));
properties.insert("agent_id".into(), json!("agt_a"));
properties.insert("command".into(), json!("status"));
HandoffEvent {
event: name.into(),
properties,
}
}
#[test]
fn a_batch_of_captured_events_is_accepted() {
let batch = HandoffBatch {
events: vec![event("command_invoked"), event("$create_alias")],
};
assert_eq!(batch.validate(), Ok(()));
}
#[test]
fn anything_but_a_telemetry_event_shape_is_refused() {
let empty = HandoffBatch { events: vec![] };
assert!(empty.validate().is_err(), "an empty hand-off");
let too_many = HandoffBatch {
events: (0..=MAX_EVENTS).map(|_| event("command_invoked")).collect(),
};
assert!(too_many.validate().is_err(), "more than MAX_EVENTS");
for name in [
"",
"Command",
"command-invoked",
"$",
"a b",
&"x".repeat(65),
] {
let batch = HandoffBatch {
events: vec![event(name)],
};
assert!(batch.validate().is_err(), "name {name:?}");
}
for key in ["distinct_id", "agent_id"] {
let mut e = event("command_invoked");
e.properties.remove(key);
let batch = HandoffBatch { events: vec![e] };
assert!(batch.validate().is_err(), "without {key}");
}
}
#[test]
fn unknown_fields_are_refused_at_parse() {
let extra_top = json!({ "events": [], "api_key": "phc_x" });
assert!(serde_json::from_value::<HandoffBatch>(extra_top).is_err());
let extra_event = json!({ "events": [{ "event": "x", "properties": {}, "timestamp": 1 }] });
assert!(serde_json::from_value::<HandoffBatch>(extra_event).is_err());
}
#[test]
fn no_daemon_drops_the_events_inside_the_bound() {
let dir = tempfile::tempdir().expect("tempdir");
std::fs::write(dir.path().join("daemon.token"), "t0ken").expect("token");
let port = std::net::TcpListener::bind("127.0.0.1:0")
.expect("bind")
.local_addr()
.expect("addr")
.port();
let events = vec![QueuedEvent {
name: "command_invoked".into(),
properties: event("command_invoked").properties,
}];
let started = std::time::Instant::now();
assert!(!send(events, dir.path(), port));
assert!(
started.elapsed() < TOTAL_TIMEOUT * 3,
"a refused hand-off took {:?}",
started.elapsed()
);
}
#[test]
fn no_pid_file_hands_nothing_off() {
use super::super::client::{start, ClientConfig};
use super::super::consent::{ConsentState, DecidedBy, Resolved};
let dir = tempfile::tempdir().expect("tempdir");
let handle = start(ClientConfig {
resolved: Resolved {
state: ConsentState::Enabled,
decided_by: DecidedBy::ConfigFile,
},
super_props: super::super::super_props::SuperProps::new("agt_a".into(), false),
debug_stderr: false,
baked_key_present: true,
egress: crate::egress::EgressConfig::direct(),
host: "http://127.0.0.1:9".into(),
});
handle.capture(super::super::events::Event::command_invoked(
"status", None, 0, 1,
));
hand_off(&handle, dir.path());
assert_eq!(
handle.take_unsent(Duration::from_secs(2)).len(),
1,
"the batch was never drained"
);
}
#[test]
fn no_token_sends_nothing() {
let dir = tempfile::tempdir().expect("tempdir");
let events = vec![QueuedEvent {
name: "command_invoked".into(),
properties: event("command_invoked").properties,
}];
assert!(!send(events, dir.path(), 9));
}
}