refactor(portal): move telemetry ingest into module
This commit is contained in:
@@ -34,6 +34,7 @@ mod risk_narrative;
|
|||||||
mod role_access;
|
mod role_access;
|
||||||
mod snapshot_cache;
|
mod snapshot_cache;
|
||||||
mod static_assets;
|
mod static_assets;
|
||||||
|
mod telemetry_ingest;
|
||||||
mod workforce_kpi_explain;
|
mod workforce_kpi_explain;
|
||||||
|
|
||||||
use api_contracts::api_contract_summary;
|
use api_contracts::api_contract_summary;
|
||||||
@@ -51,8 +52,8 @@ use path_query::{
|
|||||||
use portal_roles::PortalRole;
|
use portal_roles::PortalRole;
|
||||||
use production::{
|
use production::{
|
||||||
build_healthz, build_readyz, build_version, is_limited_api_route, mark_request_started,
|
build_healthz, build_readyz, build_version, is_limited_api_route, mark_request_started,
|
||||||
record_ingestion_accepted, record_ingestion_rejected, record_report_generated,
|
record_report_generated, render_prometheus_metrics, validate_api_query_limits,
|
||||||
render_prometheus_metrics, validate_api_query_limits, validate_portal_config,
|
validate_portal_config,
|
||||||
};
|
};
|
||||||
use readiness_api::{readiness_bundle, readiness_latest, readiness_verify};
|
use readiness_api::{readiness_bundle, readiness_latest, readiness_verify};
|
||||||
use risk_narrative::{
|
use risk_narrative::{
|
||||||
@@ -65,6 +66,12 @@ use snapshot_cache::{
|
|||||||
use static_assets::{
|
use static_assets::{
|
||||||
API_CONTRACT_OPENAPI, API_CONTRACT_TYPESCRIPT, APP_CSS, APP_JS, ARCHITECTURE_HTML, INDEX_HTML,
|
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};
|
use workforce_kpi_explain::{KpiExplainQuery, build_workforce_kpi_explain};
|
||||||
|
|
||||||
const UEBA_BASELINE_MIN_SAMPLES: usize = 3;
|
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<String> {
|
pub(crate) fn read_limited_body(request: &mut Request, limit: u64) -> Result<String> {
|
||||||
let mut body = String::new();
|
let mut body = String::new();
|
||||||
request
|
request
|
||||||
.as_reader()
|
.as_reader()
|
||||||
@@ -8570,12 +8577,12 @@ fn read_limited_body(request: &mut Request, limit: u64) -> Result<String> {
|
|||||||
Ok(body)
|
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()
|
err.to_string()
|
||||||
.contains("request payload exceeds configured limit")
|
.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(
|
respond_json_status(
|
||||||
request,
|
request,
|
||||||
StatusCode(413),
|
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<Value> {
|
|
||||||
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<IncidentActionResponse> {
|
fn apply_incident_action(args: &Cli, actor: &str, body: &str) -> Result<IncidentActionResponse> {
|
||||||
let action: IncidentActionRequest =
|
let action: IncidentActionRequest =
|
||||||
serde_json::from_str(body).map_err(|err| anyhow!("invalid incident action JSON: {err}"))?;
|
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())
|
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<String> {
|
|
||||||
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 {
|
fn request_actor(request: &Request) -> String {
|
||||||
for name in ["X-Remote-User", "X-Gateway-User", "Remote-User"] {
|
for name in ["X-Remote-User", "X-Gateway-User", "Remote-User"] {
|
||||||
if let Some(value) = request
|
if let Some(value) = request
|
||||||
|
|||||||
@@ -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<String> {
|
||||||
|
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<Value> {
|
||||||
|
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(())
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user