From 66b0b6cffb237da9716cbbcb87db323e5b9af66d Mon Sep 17 00:00:00 2001 From: igor04091968 Date: Mon, 15 Jun 2026 14:24:10 +0300 Subject: [PATCH] refactor(portal): move telemetry ingest into module --- adk-rust/crates/detmir-portal/src/main.rs | 182 ++--------------- .../detmir-portal/src/telemetry_ingest.rs | 185 ++++++++++++++++++ 2 files changed, 197 insertions(+), 170 deletions(-) create mode 100644 adk-rust/crates/detmir-portal/src/telemetry_ingest.rs diff --git a/adk-rust/crates/detmir-portal/src/main.rs b/adk-rust/crates/detmir-portal/src/main.rs index 7c046c8..fd2a020 100644 --- a/adk-rust/crates/detmir-portal/src/main.rs +++ b/adk-rust/crates/detmir-portal/src/main.rs @@ -34,6 +34,7 @@ mod risk_narrative; mod role_access; mod snapshot_cache; mod static_assets; +mod telemetry_ingest; mod workforce_kpi_explain; use api_contracts::api_contract_summary; @@ -51,8 +52,8 @@ use path_query::{ use portal_roles::PortalRole; use production::{ build_healthz, build_readyz, build_version, is_limited_api_route, mark_request_started, - record_ingestion_accepted, record_ingestion_rejected, record_report_generated, - render_prometheus_metrics, validate_api_query_limits, validate_portal_config, + record_report_generated, render_prometheus_metrics, validate_api_query_limits, + validate_portal_config, }; use readiness_api::{readiness_bundle, readiness_latest, readiness_verify}; use risk_narrative::{ @@ -65,6 +66,12 @@ use snapshot_cache::{ use static_assets::{ API_CONTRACT_OPENAPI, API_CONTRACT_TYPESCRIPT, APP_CSS, APP_JS, ARCHITECTURE_HTML, INDEX_HTML, }; +#[cfg(test)] +#[allow(unused_imports)] +pub(crate) use telemetry_ingest::{ + apply_telemetry_ingest, telemetry_authorized, validate_telemetry_payload, +}; +pub(crate) use telemetry_ingest::{bearer_token, constant_time_eq, handle_telemetry_ingest}; use workforce_kpi_explain::{KpiExplainQuery, build_workforce_kpi_explain}; const UEBA_BASELINE_MIN_SAMPLES: usize = 3; @@ -8558,7 +8565,7 @@ fn handle_incident_action(mut request: Request, args: &Cli) -> Result<()> { } } -fn read_limited_body(request: &mut Request, limit: u64) -> Result { +pub(crate) fn read_limited_body(request: &mut Request, limit: u64) -> Result { let mut body = String::new(); request .as_reader() @@ -8570,12 +8577,12 @@ fn read_limited_body(request: &mut Request, limit: u64) -> Result { Ok(body) } -fn is_payload_too_large(err: &anyhow::Error) -> bool { +pub(crate) fn is_payload_too_large(err: &anyhow::Error) -> bool { err.to_string() .contains("request payload exceeds configured limit") } -fn respond_payload_too_large(request: Request) -> Result<()> { +pub(crate) fn respond_payload_too_large(request: Request) -> Result<()> { respond_json_status( request, StatusCode(413), @@ -8734,130 +8741,6 @@ fn handle_case_details( } } -fn handle_telemetry_ingest(mut request: Request, args: &Cli) -> Result<()> { - if !telemetry_authorized(&request, args) { - record_ingestion_rejected(); - return respond_json_status( - request, - StatusCode(401), - &json!({ - "ok": false, - "error": "telemetry api key is missing or invalid" - }), - ); - } - let telemetry_limit = args.max_request_body_bytes.min(1024 * 1024); - let body = match read_limited_body(&mut request, telemetry_limit) { - Ok(body) => body, - Err(err) if is_payload_too_large(&err) => { - record_ingestion_rejected(); - return respond_payload_too_large(request); - } - Err(err) => return Err(err), - }; - let response = apply_telemetry_ingest(args, &body); - match response { - Ok(response) => { - record_ingestion_accepted(); - respond_json(request, &response) - } - Err(err) => { - record_ingestion_rejected(); - respond_json_status( - request, - StatusCode(400), - &json!({ - "ok": false, - "error": err.to_string() - }), - ) - } - } -} - -fn apply_telemetry_ingest(args: &Cli, body: &str) -> Result { - let payload: Value = - serde_json::from_str(body).map_err(|err| anyhow!("invalid telemetry JSON: {err}"))?; - validate_telemetry_payload(&payload)?; - let received_at_utc = now(); - let envelope = json!({ - "received_at_utc": received_at_utc, - "prototype": true, - "record": payload, - }); - if let Some(parent) = args.telemetry_store_path.parent() { - fs::create_dir_all(parent).with_context(|| format!("create {}", parent.display()))?; - } - let mut file = OpenOptions::new() - .create(true) - .append(true) - .open(&args.telemetry_store_path) - .with_context(|| format!("open {}", args.telemetry_store_path.display()))?; - writeln!(file, "{}", serde_json::to_string(&envelope)?) - .with_context(|| format!("append {}", args.telemetry_store_path.display()))?; - Ok(json!({ - "ok": true, - "prototype": true, - "stored": "file-backed-jsonl", - "received_at_utc": received_at_utc, - })) -} - -fn validate_telemetry_payload(payload: &Value) -> Result<()> { - let Some(object) = payload.as_object() else { - return Err(anyhow!("telemetry payload must be a JSON object")); - }; - for field in [ - "agent_id", - "hostname", - "os_name", - "os_version", - "platform", - "username", - "timestamp", - "uptime_seconds", - "cpu_usage_percent", - "memory_total", - "memory_used", - "active_sessions", - "rdp_sessions", - "ssh_sessions", - "processes", - "network_interfaces", - "network_connections", - "workforce_activity", - "security_events", - "collector_version", - ] { - if !object.contains_key(field) { - return Err(anyhow!("telemetry field is missing: {field}")); - } - } - for field in [ - "active_sessions", - "rdp_sessions", - "ssh_sessions", - "processes", - "network_interfaces", - "network_connections", - "security_events", - ] { - if payload.get(field).and_then(Value::as_array).is_none() { - return Err(anyhow!("telemetry field must be an array: {field}")); - } - } - if payload - .get("workforce_activity") - .and_then(Value::as_object) - .is_none() - { - return Err(anyhow!( - "telemetry field must be an object: workforce_activity" - )); - } - Ok(()) -} - fn apply_incident_action(args: &Cli, actor: &str, body: &str) -> Result { let action: IncidentActionRequest = serde_json::from_str(body).map_err(|err| anyhow!("invalid incident action JSON: {err}"))?; @@ -9968,47 +9851,6 @@ fn upload_authorized(request: &Request, args: &Cli) -> bool { constant_time_eq(actual.as_bytes(), expected.as_bytes()) } -fn telemetry_authorized(request: &Request, args: &Cli) -> bool { - let expected = args.telemetry_api_key.trim(); - if expected.is_empty() || expected == "change-me" { - return false; - } - let actual = request - .headers() - .iter() - .find(|header| header.field.equiv("x-api-key")) - .map(|header| header.value.as_str().trim().to_string()) - .or_else(|| bearer_token(request)); - actual - .as_deref() - .filter(|value| !value.is_empty()) - .map(|value| constant_time_eq(value.as_bytes(), expected.as_bytes())) - .unwrap_or(false) -} - -fn bearer_token(request: &Request) -> Option { - request - .headers() - .iter() - .find(|header| header.field.equiv("Authorization")) - .map(|header| header.value.as_str().trim()) - .and_then(|value| value.strip_prefix("Bearer ")) - .map(str::trim) - .filter(|value| !value.is_empty()) - .map(ToString::to_string) -} - -fn constant_time_eq(left: &[u8], right: &[u8]) -> bool { - if left.len() != right.len() { - return false; - } - let mut diff = 0_u8; - for (a, b) in left.iter().zip(right.iter()) { - diff |= a ^ b; - } - diff == 0 -} - fn request_actor(request: &Request) -> String { for name in ["X-Remote-User", "X-Gateway-User", "Remote-User"] { if let Some(value) = request diff --git a/adk-rust/crates/detmir-portal/src/telemetry_ingest.rs b/adk-rust/crates/detmir-portal/src/telemetry_ingest.rs new file mode 100644 index 0000000..6cd27a3 --- /dev/null +++ b/adk-rust/crates/detmir-portal/src/telemetry_ingest.rs @@ -0,0 +1,185 @@ +//! Telemetry ingest API for Rust endpoint/agent diagnostics. +//! +//! CONTRACT: this module owns `/api/telemetry` authentication, request-body +//! validation and JSONL append semantics. Keep accepted fields, status codes, +//! metrics counters and response shape stable unless telemetry contracts are +//! updated in the same PR. + +use std::fs; +use std::fs::OpenOptions; +use std::io::Write; + +use anyhow::{Context, Result, anyhow}; +use serde_json::{Value, json}; +use tiny_http::{Request, StatusCode}; + +use crate::production::{record_ingestion_accepted, record_ingestion_rejected}; +use crate::{ + Cli, is_payload_too_large, now, read_limited_body, respond_json, respond_json_status, + respond_payload_too_large, +}; + +pub(crate) fn telemetry_authorized(request: &Request, args: &Cli) -> bool { + let expected = args.telemetry_api_key.trim(); + if expected.is_empty() || expected == "change-me" { + return false; + } + let actual = request + .headers() + .iter() + .find(|header| header.field.equiv("x-api-key")) + .map(|header| header.value.as_str().trim().to_string()) + .or_else(|| bearer_token(request)); + actual + .as_deref() + .filter(|value| !value.is_empty()) + .map(|value| constant_time_eq(value.as_bytes(), expected.as_bytes())) + .unwrap_or(false) +} + +pub(crate) fn bearer_token(request: &Request) -> Option { + request + .headers() + .iter() + .find(|header| header.field.equiv("Authorization")) + .map(|header| header.value.as_str().trim()) + .and_then(|value| value.strip_prefix("Bearer ")) + .map(str::trim) + .filter(|value| !value.is_empty()) + .map(ToString::to_string) +} + +pub(crate) fn constant_time_eq(left: &[u8], right: &[u8]) -> bool { + if left.len() != right.len() { + return false; + } + let mut diff = 0_u8; + for (a, b) in left.iter().zip(right.iter()) { + diff |= a ^ b; + } + diff == 0 +} + +pub(crate) fn handle_telemetry_ingest(mut request: Request, args: &Cli) -> Result<()> { + if !telemetry_authorized(&request, args) { + record_ingestion_rejected(); + return respond_json_status( + request, + StatusCode(401), + &json!({ + "ok": false, + "error": "telemetry api key is missing or invalid" + }), + ); + } + let telemetry_limit = args.max_request_body_bytes.min(1024 * 1024); + let body = match read_limited_body(&mut request, telemetry_limit) { + Ok(body) => body, + Err(err) if is_payload_too_large(&err) => { + record_ingestion_rejected(); + return respond_payload_too_large(request); + } + Err(err) => return Err(err), + }; + let response = apply_telemetry_ingest(args, &body); + match response { + Ok(response) => { + record_ingestion_accepted(); + respond_json(request, &response) + } + Err(err) => { + record_ingestion_rejected(); + respond_json_status( + request, + StatusCode(400), + &json!({ + "ok": false, + "error": err.to_string() + }), + ) + } + } +} + +pub(crate) fn apply_telemetry_ingest(args: &Cli, body: &str) -> Result { + let payload: Value = + serde_json::from_str(body).map_err(|err| anyhow!("invalid telemetry JSON: {err}"))?; + validate_telemetry_payload(&payload)?; + let received_at_utc = now(); + let envelope = json!({ + "received_at_utc": received_at_utc, + "prototype": true, + "record": payload, + }); + if let Some(parent) = args.telemetry_store_path.parent() { + fs::create_dir_all(parent).with_context(|| format!("create {}", parent.display()))?; + } + let mut file = OpenOptions::new() + .create(true) + .append(true) + .open(&args.telemetry_store_path) + .with_context(|| format!("open {}", args.telemetry_store_path.display()))?; + writeln!(file, "{}", serde_json::to_string(&envelope)?) + .with_context(|| format!("append {}", args.telemetry_store_path.display()))?; + Ok(json!({ + "ok": true, + "prototype": true, + "stored": "file-backed-jsonl", + "received_at_utc": received_at_utc, + })) +} + +pub(crate) fn validate_telemetry_payload(payload: &Value) -> Result<()> { + let Some(object) = payload.as_object() else { + return Err(anyhow!("telemetry payload must be a JSON object")); + }; + for field in [ + "agent_id", + "hostname", + "os_name", + "os_version", + "platform", + "username", + "timestamp", + "uptime_seconds", + "cpu_usage_percent", + "memory_total", + "memory_used", + "active_sessions", + "rdp_sessions", + "ssh_sessions", + "processes", + "network_interfaces", + "network_connections", + "workforce_activity", + "security_events", + "collector_version", + ] { + if !object.contains_key(field) { + return Err(anyhow!("telemetry field is missing: {field}")); + } + } + for field in [ + "active_sessions", + "rdp_sessions", + "ssh_sessions", + "processes", + "network_interfaces", + "network_connections", + "security_events", + ] { + if payload.get(field).and_then(Value::as_array).is_none() { + return Err(anyhow!("telemetry field must be an array: {field}")); + } + } + if payload + .get("workforce_activity") + .and_then(Value::as_object) + .is_none() + { + return Err(anyhow!( + "telemetry field must be an object: workforce_activity" + )); + } + Ok(()) +}