From 21f9826e68d31602ef2185f87ab3c573e1ed2360 Mon Sep 17 00:00:00 2001 From: Aleksei Fedorov Date: Fri, 31 Jul 2026 07:49:35 +0000 Subject: [PATCH] ML collectors, attribution, rules, sequence deep --- kernel_ai/ml/attribution/__init__.py | 5 + kernel_ai/ml/attribution/attack_map.py | 170 +++++++++++++++++ kernel_ai/ml/attribution/classifier.py | 12 ++ kernel_ai/ml/attribution/enrich.py | 61 ++++++ kernel_ai/ml/attribution/sigma_engine.py | 89 +++++++++ kernel_ai/ml/collectors/__init__.py | 11 ++ kernel_ai/ml/collectors/base.py | 180 ++++++++++++++++++ kernel_ai/ml/collectors/socket_source.py | 104 ++++++++++ kernel_ai/ml/collectors/stream_e2e.py | 131 +++++++++++++ kernel_ai/ml/rules/sigma/privesc_euid.json | 16 ++ .../ml/rules/sigma/reverse_shell_lineage.json | 18 ++ kernel_ai/ml/rules/sigma/shell_from_web.json | 18 ++ kernel_ai/ml/sequence_deep/__init__.py | 9 + kernel_ai/ml/sequence_deep/__main__.py | 65 +++++++ kernel_ai/ml/sequence_deep/encode.py | 66 +++++++ kernel_ai/ml/sequence_deep/lstm.py | 47 +++++ kernel_ai/ml/sequence_deep/markov.py | 88 +++++++++ kernel_ai/ml/sequence_deep/scorer.py | 118 ++++++++++++ kernel_ai/ml/sequence_deep/train_markov.py | 159 ++++++++++++++++ tests/test_ml_collectors.py | 32 ++++ tests/test_ml_polygon.py | 38 ++++ tests/test_ml_stage5.py | 74 +++++++ tests/test_ml_stage7.py | 66 +++++++ tests/test_ml_stage8.py | 119 ++++++++++++ 24 files changed, 1696 insertions(+) create mode 100644 kernel_ai/ml/attribution/__init__.py create mode 100644 kernel_ai/ml/attribution/attack_map.py create mode 100644 kernel_ai/ml/attribution/classifier.py create mode 100644 kernel_ai/ml/attribution/enrich.py create mode 100644 kernel_ai/ml/attribution/sigma_engine.py create mode 100644 kernel_ai/ml/collectors/__init__.py create mode 100644 kernel_ai/ml/collectors/base.py create mode 100644 kernel_ai/ml/collectors/socket_source.py create mode 100644 kernel_ai/ml/collectors/stream_e2e.py create mode 100644 kernel_ai/ml/rules/sigma/privesc_euid.json create mode 100644 kernel_ai/ml/rules/sigma/reverse_shell_lineage.json create mode 100644 kernel_ai/ml/rules/sigma/shell_from_web.json create mode 100644 kernel_ai/ml/sequence_deep/__init__.py create mode 100644 kernel_ai/ml/sequence_deep/__main__.py create mode 100644 kernel_ai/ml/sequence_deep/encode.py create mode 100644 kernel_ai/ml/sequence_deep/lstm.py create mode 100644 kernel_ai/ml/sequence_deep/markov.py create mode 100644 kernel_ai/ml/sequence_deep/scorer.py create mode 100644 kernel_ai/ml/sequence_deep/train_markov.py create mode 100644 tests/test_ml_collectors.py create mode 100644 tests/test_ml_polygon.py create mode 100644 tests/test_ml_stage5.py create mode 100644 tests/test_ml_stage7.py create mode 100644 tests/test_ml_stage8.py diff --git a/kernel_ai/ml/attribution/__init__.py b/kernel_ai/ml/attribution/__init__.py new file mode 100644 index 0000000..79ed3dc --- /dev/null +++ b/kernel_ai/ml/attribution/__init__.py @@ -0,0 +1,5 @@ +"""Stage 7 — ATT&CK / Sigma-lite attribution for ML anomalies.""" + +from kernel_ai.ml.attribution.enrich import enrich_anomalies + +__all__ = ["enrich_anomalies"] diff --git a/kernel_ai/ml/attribution/attack_map.py b/kernel_ai/ml/attribution/attack_map.py new file mode 100644 index 0000000..8c58146 --- /dev/null +++ b/kernel_ai/ml/attribution/attack_map.py @@ -0,0 +1,170 @@ +"""MITRE ATT&CK technique catalogue + heuristic mapping from anomaly shape. + +v1 is intentionally rule-based and explainable. A supervised classifier can +plug in later (``classifier.py``) without changing the mutation contract. +""" + +from __future__ import annotations + +from dataclasses import dataclass + + +@dataclass(frozen=True) +class Technique: + mitre: str + family: str + name: str + color: str # UI accent hint + cves: tuple[str, ...] = () + + +# Compact catalogue of families we can currently speak about from Stages 1–5. +TECHNIQUES: dict[str, Technique] = { + "T1059": Technique("T1059", "execution", "Command and Scripting Interpreter", "#e8a54b"), + "T1204": Technique("T1204", "execution", "User Execution", "#e8a54b"), + "T1498": Technique("T1498", "impact", "Network Denial of Service", "#e0564e"), + "T1499": Technique("T1499", "impact", "Endpoint Denial of Service", "#e0564e"), + "T1496": Technique("T1496", "impact", "Resource Hijacking", "#c9a6ff"), + "T1071": Technique("T1071", "command_and_control", "Application Layer Protocol", "#67c8e0"), + "T1046": Technique("T1046", "discovery", "Network Service Discovery", "#8ff0d2"), + "T1068": Technique("T1068", "privilege_escalation", "Exploitation for Privilege Escalation", "#e0564e"), + "T1548": Technique("T1548", "privilege_escalation", "Abuse Elevation Control Mechanism", "#e0564e"), + "T1083": Technique("T1083", "discovery", "File and Directory Discovery", "#b8c7da"), + "T1106": Technique("T1106", "execution", "Native API", "#f0c48a"), + "T1055": Technique("T1055", "defense_evasion", "Process Injection", "#c9a6ff"), +} + + +# child_comm tokens that suggest scripting / shells / C2 helpers +_EXEC_CHILDREN = frozenset( + { + "bash", "sh", "dash", "zsh", "python", "python3", "perl", "ruby", "php", + "node", "nc", "ncat", "netcat", "socat", "curl", "wget", "busybox", + } +) +_SCAN_HINTS = frozenset({"nmap", "masscan", "zmap", "nikto"}) +_MINER_HINTS = frozenset({"xmrig", "minerd", "cpuminer", "ethminer"}) + + +def _technique_payload(tech: Technique, *, confidence: float, source: str, why: str) -> dict: + return { + "family": tech.family, + "mitre": tech.mitre, + "name": tech.name, + "label_confidence": round(max(0.0, min(1.0, confidence)), 3), + "cve": list(tech.cves), + "source": source, + "color": tech.color, + "why": why, + } + + +def map_anomaly(anomaly: dict) -> dict | None: + """Return an ``attack`` dict or None if we should leave it unattributed.""" + source = str(anomaly.get("source") or "") + feature = str(anomaly.get("feature") or "") + atype = str(anomaly.get("type") or "") + message = str(anomaly.get("message") or "").lower() + meta = anomaly.get("meta") or {} + kind = str(meta.get("kind") or "") + + # --- Stage 5 lineage / privesc --- + if source == "stage5_process" or feature.startswith("lineage:") or atype.startswith("lineage:"): + child = str(meta.get("comm") or "") + parent = str(meta.get("parent_comm") or "") + edge = f"{parent}→{child}".lower() + child_l = child.lower() + if child_l in _MINER_HINTS or any(h in edge for h in _MINER_HINTS): + return _technique_payload( + TECHNIQUES["T1496"], + confidence=0.72, + source="heuristic", + why=f"lineage suggests miner binary ({parent}→{child})", + ) + if child_l in _SCAN_HINTS: + return _technique_payload( + TECHNIQUES["T1046"], + confidence=0.7, + source="heuristic", + why=f"lineage suggests scanner ({parent}→{child})", + ) + if child_l in _EXEC_CHILDREN or kind == "lineage": + conf = 0.66 if child_l in _EXEC_CHILDREN else 0.45 + return _technique_payload( + TECHNIQUES["T1059"], + confidence=conf, + source="heuristic", + why=f"unusual process lineage {parent}→{child}", + ) + + if kind == "privesc" or "euid_root" in atype or feature.startswith("privesc:"): + return _technique_payload( + TECHNIQUES["T1548"], + confidence=0.75, + source="heuristic", + why="effective uid 0 with non-root real uid", + ) + + # --- Stage 4 / Stage 8 sequence --- + if ( + source in ("stage4_sequence", "stage8_sequence") + or feature in ("syscall_seq", "syscall_seq_deep") + or atype in ("syscall_sequence", "syscall_sequence_deep") + ): + why = ( + "low-likelihood syscall order (deep sequence model)" + if source == "stage8_sequence" or feature == "syscall_seq_deep" + else "novel syscall sequencing (STIDE mismatch)" + ) + conf = 0.58 if source == "stage8_sequence" else 0.55 + return _technique_payload( + TECHNIQUES["T1106"], + confidence=conf, + source="heuristic", + why=why, + ) + + # --- Stage 1 / 2 host features --- + feat = feature.lower() + if any(k in feat for k in ("tcp_retrans", "net_softirq", "tcp_inseg", "tcp_outseg")): + return _technique_payload( + TECHNIQUES["T1498"], + confidence=0.4, + source="heuristic", + why=f"network-path pressure on {feature}", + ) + if any(k in feat for k in ("ctxt_per_sec", "procs_running", "proc_count", "run_queue", "load1", "cpu_busy")): + # Resource pressure — could be DoS or miner; stay conservative. + if "miner" in message or "xmrig" in message: + tech = TECHNIQUES["T1496"] + conf = 0.6 + why = "host CPU/sched pressure with miner hint" + else: + tech = TECHNIQUES["T1499"] + conf = 0.38 + why = f"host sched/CPU pressure on {feature}" + return _technique_payload(tech, confidence=conf, source="heuristic", why=why) + if any(k in feat for k in ("pgfault", "pgmajfault", "pgscan", "swap_io", "psi_mem")): + return _technique_payload( + TECHNIQUES["T1499"], + confidence=0.36, + source="heuristic", + why=f"memory pressure on {feature}", + ) + if "hardirq" in feat or "block_softirq" in feat: + return _technique_payload( + TECHNIQUES["T1499"], + confidence=0.34, + source="heuristic", + why=f"IRQ/block pressure on {feature}", + ) + + if source == "stage2_isoforest": + return _technique_payload( + TECHNIQUES["T1499"], + confidence=0.3, + source="heuristic", + why="IsolationForest unusual host state (family-level guess)", + ) + + return None diff --git a/kernel_ai/ml/attribution/classifier.py b/kernel_ai/ml/attribution/classifier.py new file mode 100644 index 0000000..bdfbee2 --- /dev/null +++ b/kernel_ai/ml/attribution/classifier.py @@ -0,0 +1,12 @@ +"""Placeholder for a future supervised ATT&CK classifier (GB / sklearn). + +v1 Stage 7 uses heuristics + Sigma-lite only. This module exists so the +roadmap file layout stays stable when polygon-labeled training lands. +""" + +from __future__ import annotations + + +def predict_attack(_anomaly: dict) -> dict | None: + """Return attack dict or None. Untrained stub always returns None.""" + return None diff --git a/kernel_ai/ml/attribution/enrich.py b/kernel_ai/ml/attribution/enrich.py new file mode 100644 index 0000000..d965400 --- /dev/null +++ b/kernel_ai/ml/attribution/enrich.py @@ -0,0 +1,61 @@ +"""Enrich ML anomaly records with Stage 7 ``attack`` attribution.""" + +from __future__ import annotations + +import logging + +from kernel_ai.ml.attribution import classifier +from kernel_ai.ml.attribution.attack_map import map_anomaly +from kernel_ai.ml.attribution.sigma_engine import load_rules, match_anomaly + +logger = logging.getLogger("kernel_ai.ml.attribution") + +# Below this confidence we keep the anomaly but mark family as unknown-ish +# by omitting attack (UI stays uncolored). Tunable via enrich() arg. +_DEFAULT_MIN_CONF = 0.35 + + +def enrich_anomaly( + anomaly: dict, + *, + rules: list[dict] | None = None, + min_confidence: float = _DEFAULT_MIN_CONF, +) -> dict: + """Attach ``attack`` + mirror into ``meta.attack`` (back-compat for DNA).""" + if not isinstance(anomaly, dict): + return anomaly + if anomaly.get("attack"): + return anomaly + + # Precedence: Sigma-lite (high precision) > heuristic map > ML classifier stub. + attack = match_anomaly(anomaly, rules=rules) + if attack is None: + attack = map_anomaly(anomaly) + if attack is None: + attack = classifier.predict_attack(anomaly) + if attack is None: + return anomaly + if float(attack.get("label_confidence") or 0.0) < min_confidence: + return anomaly + + anomaly = dict(anomaly) + anomaly["attack"] = attack + meta = dict(anomaly.get("meta") or {}) + meta["attack"] = attack + anomaly["meta"] = meta + # Optional: prefix message once for operators curling the API. + mitre = attack.get("mitre") + if mitre and mitre not in str(anomaly.get("message") or ""): + anomaly["message"] = f"[{mitre}] {anomaly.get('message') or ''}".strip() + return anomaly + + +def enrich_anomalies(anomalies: list[dict], *, min_confidence: float = _DEFAULT_MIN_CONF) -> list[dict]: + if not anomalies: + return [] + try: + rules = load_rules() + except Exception as exc: # noqa: BLE001 + logger.warning("sigma rules load failed: %s", exc) + rules = [] + return [enrich_anomaly(a, rules=rules, min_confidence=min_confidence) for a in anomalies] diff --git a/kernel_ai/ml/attribution/sigma_engine.py b/kernel_ai/ml/attribution/sigma_engine.py new file mode 100644 index 0000000..3d2d00a --- /dev/null +++ b/kernel_ai/ml/attribution/sigma_engine.py @@ -0,0 +1,89 @@ +"""Minimal Sigma-lite matcher (JSON rules, no PyYAML dependency). + +Full Sigma is Stage 7 aspirational; this covers a few high-precision patterns +we can express from ML anomaly records + Stage 5 meta (comm/lineage). +""" + +from __future__ import annotations + +import json +import logging +from pathlib import Path + +logger = logging.getLogger("kernel_ai.ml.attribution.sigma") + +_RULES_DIR = Path(__file__).resolve().parents[1] / "rules" / "sigma" + + +def load_rules(directory: Path | None = None) -> list[dict]: + root = directory or _RULES_DIR + rules: list[dict] = [] + if not root.is_dir(): + return rules + for path in sorted(root.glob("*.json")): + try: + data = json.loads(path.read_text(encoding="utf-8")) + except (OSError, json.JSONDecodeError) as exc: + logger.warning("skip sigma rule %s: %s", path.name, exc) + continue + if isinstance(data, dict): + data.setdefault("id", path.stem) + rules.append(data) + elif isinstance(data, list): + for i, item in enumerate(data): + if isinstance(item, dict): + item.setdefault("id", f"{path.stem}_{i}") + rules.append(item) + return rules + + +def match_anomaly(anomaly: dict, rules: list[dict] | None = None) -> dict | None: + """Return attack dict from the first matching rule, else None.""" + rules = rules if rules is not None else load_rules() + meta = anomaly.get("meta") or {} + hay = { + "source": str(anomaly.get("source") or ""), + "feature": str(anomaly.get("feature") or ""), + "type": str(anomaly.get("type") or ""), + "message": str(anomaly.get("message") or ""), + "comm": str(meta.get("comm") or ""), + "parent_comm": str(meta.get("parent_comm") or ""), + "kind": str(meta.get("kind") or ""), + } + edge = f"{hay['parent_comm']}->{hay['comm']}".lower() + + for rule in rules: + when = rule.get("when") or {} + ok = True + for key, expected in when.items(): + if key == "child_in": + ok = hay["comm"].lower() in {str(x).lower() for x in expected} + elif key == "parent_in": + ok = hay["parent_comm"].lower() in {str(x).lower() for x in expected} + elif key == "edge_contains": + ok = any(str(x).lower() in edge for x in expected) + elif key == "feature_contains": + ok = any(str(x).lower() in hay["feature"].lower() for x in expected) + elif key == "source_in": + ok = hay["source"] in set(expected) + elif key == "kind_in": + ok = hay["kind"] in set(expected) + else: + ok = True + if not ok: + break + if not ok: + continue + attack = rule.get("attack") or {} + return { + "family": attack.get("family", "unknown"), + "mitre": attack.get("mitre"), + "name": attack.get("name") or rule.get("title") or rule.get("id"), + "label_confidence": float(attack.get("confidence", 0.85)), + "cve": list(attack.get("cve") or []), + "source": "sigma", + "color": attack.get("color", "#e8a54b"), + "why": rule.get("description") or rule.get("title") or rule.get("id"), + "rule_id": rule.get("id"), + } + return None diff --git a/kernel_ai/ml/collectors/__init__.py b/kernel_ai/ml/collectors/__init__.py new file mode 100644 index 0000000..4ba75e1 --- /dev/null +++ b/kernel_ai/ml/collectors/__init__.py @@ -0,0 +1,11 @@ +"""Syscall event sources for Stage 4 / Stage 6. + +Privileged collection lives outside the ML worker (see +``docs/ML_STAGE6_L2_COLLECTOR.md``). The worker only *drains* a normalized +event stream and feeds :class:`kernel_ai.ml.sequence.NgramTracker`. +""" + +from kernel_ai.ml.collectors.base import ALLOWED_SYSCALLS, SyscallEvent +from kernel_ai.ml.collectors.socket_source import SocketSyscallSource + +__all__ = ["ALLOWED_SYSCALLS", "SyscallEvent", "SocketSyscallSource"] diff --git a/kernel_ai/ml/collectors/base.py b/kernel_ai/ml/collectors/base.py new file mode 100644 index 0000000..6065ec7 --- /dev/null +++ b/kernel_ai/ml/collectors/base.py @@ -0,0 +1,180 @@ +"""Shared syscall-event contract for Stage 6 L2 collectors.""" + +from __future__ import annotations + +import json +from dataclasses import asdict, dataclass +from typing import Iterable, Iterator + + +# Security-relevant allowlist for v1 (keep CPU bounded on small hosts). +ALLOWED_SYSCALLS: frozenset[str] = frozenset( + { + "execve", + "execveat", + "clone", + "clone3", + "fork", + "vfork", + "connect", + "accept", + "accept4", + "bind", + "listen", + "open", + "openat", + "creat", + "unlinkat", + "renameat", + "renameat2", + "mmap", + "mprotect", + "pkey_mprotect", + "setuid", + "setreuid", + "setresuid", + "setgid", + "setregid", + "setresgid", + "ptrace", + "process_vm_writev", + "memfd_create", + "userfaultfd", + } +) + +# Linux x86_64 syscall numbers for the allowlist (audit logs emit numbers). +ALLOWED_SYSCALL_NR: frozenset[int] = frozenset( + { + 56, # clone + 57, # fork + 58, # vfork + 59, # execve + 322, # execveat + 435, # clone3 + 41, # socket (not scored alone; kept out — connect/accept matter more) + 42, # connect + 43, # accept + 49, # bind + 50, # listen + 288, # accept4 + 2, # open + 257, # openat + 85, # creat + 263, # unlinkat + 264, # renameat + 316, # renameat2 + 9, # mmap + 10, # mprotect + 330, # pkey_mprotect + 105, # setuid + 113, # setreuid + 117, # setresuid + 106, # setgid + 114, # setregid + 119, # setresgid + 101, # ptrace + 310, # process_vm_writev + 319, # memfd_create + 323, # userfaultfd + } +) + +# Minimal nr → name map for allowlisted calls (collector / audit parser). +SYSCALL_NR_TO_NAME: dict[int, str] = { + 56: "clone", + 57: "fork", + 58: "vfork", + 59: "execve", + 322: "execveat", + 435: "clone3", + 42: "connect", + 43: "accept", + 49: "bind", + 50: "listen", + 288: "accept4", + 2: "open", + 257: "openat", + 85: "creat", + 263: "unlinkat", + 264: "renameat", + 316: "renameat2", + 9: "mmap", + 10: "mprotect", + 330: "pkey_mprotect", + 105: "setuid", + 113: "setreuid", + 117: "setresuid", + 106: "setgid", + 114: "setregid", + 119: "setresgid", + 101: "ptrace", + 310: "process_vm_writev", + 319: "memfd_create", + 323: "userfaultfd", +} + + +@dataclass(frozen=True) +class SyscallEvent: + """Normalized L2 syscall event (collector → ML worker).""" + + ts: float + pid: int + uid: int + comm: str + syscall: str + + def to_json(self) -> str: + return json.dumps(asdict(self), separators=(",", ":")) + + @classmethod + def from_mapping(cls, data: dict) -> "SyscallEvent | None": + try: + syscall = str(data.get("syscall") or "").strip() + if not syscall: + return None + return cls( + ts=float(data.get("ts") or 0.0), + pid=int(data.get("pid")), + uid=int(data.get("uid") or 0), + comm=str(data.get("comm") or "?")[:64], + syscall=syscall[:64], + ) + except (TypeError, ValueError): + return None + + +def encode_events(events: Iterable[SyscallEvent]) -> bytes: + """Pack one or more events into a single datagram payload.""" + body = "\n".join(ev.to_json() for ev in events) + return (body + "\n").encode("utf-8") + + +def decode_events(payload: bytes) -> list[SyscallEvent]: + """Parse a datagram that may contain multiple NDJSON lines.""" + out: list[SyscallEvent] = [] + try: + text = payload.decode("utf-8", errors="ignore") + except Exception: + return out + for line in text.splitlines(): + line = line.strip() + if not line: + continue + try: + data = json.loads(line) + except json.JSONDecodeError: + continue + if not isinstance(data, dict): + continue + ev = SyscallEvent.from_mapping(data) + if ev is not None: + out.append(ev) + return out + + +def iter_allowed(events: Iterable[SyscallEvent]) -> Iterator[SyscallEvent]: + for ev in events: + if ev.syscall in ALLOWED_SYSCALLS or ev.syscall.startswith("sys_"): + yield ev diff --git a/kernel_ai/ml/collectors/socket_source.py b/kernel_ai/ml/collectors/socket_source.py new file mode 100644 index 0000000..8282a48 --- /dev/null +++ b/kernel_ai/ml/collectors/socket_source.py @@ -0,0 +1,104 @@ +"""Worker-side reader for the Stage 6 unix-datagram syscall stream.""" + +from __future__ import annotations + +import logging +import os +import socket +import stat +from typing import List + +from kernel_ai.ml.collectors.base import SyscallEvent, decode_events + +logger = logging.getLogger("kernel_ai.ml.collectors.socket") + +_MAX_DATAGRAM = 65535 + + +class SocketSyscallSource: + """Bind ``/run/kernel-ai/ml-syscall.sock`` and drain collector datagrams. + + Push model: the privileged collector ``sendto``s this path; the ML worker + (www-data) owns the bind. If bind fails, :meth:`drain` returns [] and + Stage 4 stays dormant without affecting Stages 1–2. + """ + + def __init__(self, path: str, *, max_events: int = 2000) -> None: + self.path = path + self.max_events = max(1, max_events) + self._sock: socket.socket | None = None + self._warned = False + + def close(self) -> None: + if self._sock is not None: + try: + self._sock.close() + except OSError: + pass + self._sock = None + try: + if os.path.exists(self.path): + os.unlink(self.path) + except OSError: + pass + + def _ensure_sock(self) -> socket.socket | None: + if self._sock is not None: + return self._sock + directory = os.path.dirname(self.path) or "." + try: + os.makedirs(directory, mode=0o775, exist_ok=True) + except OSError as exc: + if not self._warned: + logger.warning("seq socket dir %s: %s", directory, exc) + self._warned = True + return None + try: + if os.path.exists(self.path): + os.unlink(self.path) + except OSError: + pass + sock = socket.socket(socket.AF_UNIX, socket.SOCK_DGRAM) + sock.setblocking(False) + try: + sock.bind(self.path) + os.chmod(self.path, stat.S_IRUSR | stat.S_IWUSR | stat.S_IRGRP | stat.S_IWGRP) + except OSError as exc: + sock.close() + if not self._warned: + logger.warning("seq socket bind failed (%s): %s", self.path, exc) + self._warned = True + return None + # Enlarge recv buffer so short collector bursts are less likely to drop. + try: + sock.setsockopt(socket.SOL_SOCKET, socket.SO_RCVBUF, 1 << 20) + except OSError: + pass + self._sock = sock + self._warned = False + logger.info("listening on seq socket %s", self.path) + return sock + + def drain(self, max_events: int | None = None) -> List[SyscallEvent]: + """Receive pending datagrams; return up to ``max_events`` events.""" + limit = self.max_events if max_events is None else max(1, max_events) + sock = self._ensure_sock() + if sock is None: + return [] + out: list[SyscallEvent] = [] + while len(out) < limit: + try: + payload = sock.recv(_MAX_DATAGRAM) + except BlockingIOError: + break + except OSError as exc: + logger.warning("seq socket recv failed: %s — rebinding", exc) + self.close() + break + if not payload: + break + for ev in decode_events(payload): + out.append(ev) + if len(out) >= limit: + break + return out diff --git a/kernel_ai/ml/collectors/stream_e2e.py b/kernel_ai/ml/collectors/stream_e2e.py new file mode 100644 index 0000000..a39ff1e --- /dev/null +++ b/kernel_ai/ml/collectors/stream_e2e.py @@ -0,0 +1,131 @@ +"""Stage 6 local end-to-end: unix-datagram collector → NgramTracker → STIDE/Markov. + +No root / auditd required. Uses the same socket contract as PROD +(``SocketSyscallSource`` bind + collector ``sendto``). +""" + +from __future__ import annotations + +import os +import tempfile +import threading +import time +from pathlib import Path + +from kernel_ai.ml.collectors.socket_source import SocketSyscallSource +from kernel_ai.ml.sequence import NgramTracker, StideModel +from kernel_ai.ml.sequence_deep.markov import MarkovScorer + + +def _token_ngrams(tokens: list[str], n: int = 3) -> list[str]: + if len(tokens) < n: + return [] + return ["|".join(tokens[i : i + n]) for i in range(len(tokens) - n + 1)] + + +def run_stream_e2e( + *, + bursts: int = 9, + demo_every: float = 0.05, + socket_path: str | None = None, +) -> dict: + """Drive demo emitter → socket → tracker; score STIDE + Markov. + + Returns a result dict with ``pass`` True when the socket path delivered + events and Markov (trained on *normal* demo chains only) separates + normal vs mimicry windows. + """ + # Import demo helpers from the collector script module. + import importlib.util + + collector_path = Path(__file__).resolve().parents[3] / "deploy" / "ebpf" / "syscall_stream_collector.py" + spec = importlib.util.spec_from_file_location("kai_syscall_collector", collector_path) + if spec is None or spec.loader is None: + raise RuntimeError(f"cannot load collector from {collector_path}") + mod = importlib.util.module_from_spec(spec) + spec.loader.exec_module(mod) + + tmp_dir = tempfile.mkdtemp(prefix="kai-seq-") + path = socket_path or os.path.join(tmp_dir, "ml-syscall.sock") + + source = SocketSyscallSource(path, max_events=5000) + # Force bind before emitter starts. + assert source._ensure_sock() is not None, f"failed to bind {path}" + + emitter = mod.DatagramEmitter(path) + stop = threading.Event() + + def _emit() -> None: + old = mod.DEMO_EVERY + mod.DEMO_EVERY = demo_every + try: + mod.run_demo(emitter, bursts=bursts) + finally: + mod.DEMO_EVERY = old + stop.set() + + thread = threading.Thread(target=_emit, name="seq-demo-emit", daemon=True) + thread.start() + + tracker = NgramTracker(n=3, window=800) + raw_tokens: list[str] = [] + normal_tokens: list[str] = [] + deadline = time.time() + max(5.0, bursts * demo_every + 2.0) + while time.time() < deadline: + events = source.drain() + if events: + tracker.update_stream(events) + raw_tokens.extend(ev.syscall for ev in events) + # Train only on quiet host-like bursts (comm=demo). + normal_tokens.extend(ev.syscall for ev in events if ev.comm == "demo") + if stop.is_set() and not events: + time.sleep(0.05) + events = source.drain() + if events: + tracker.update_stream(events) + raw_tokens.extend(ev.syscall for ev in events) + normal_tokens.extend(ev.syscall for ev in events if ev.comm == "demo") + break + time.sleep(0.02) + + thread.join(timeout=2) + source.close() + emitter.close() + + window = tracker.recent() + markov = MarkovScorer(meta={"source": "stream_e2e"}) + for i in range(0, max(0, len(normal_tokens) - 4), 4): + chunk = normal_tokens[i : i + 8] + if len(chunk) >= 2: + markov.observe(chunk) + + stide = StideModel(n=3, ngrams=set(window), meta={"source": "stream_e2e"}) + + normal = list(mod.DEMO_NORMAL_CHAINS[0]) * 4 + mimic = list(mod.DEMO_MIMICRY_CHAIN) * 4 + n_score = (markov.score_window(normal) or {}).get("neg_avg_logprob", 0.0) + m_score = (markov.score_window(mimic) or {}).get("neg_avg_logprob", 0.0) + # STIDE on mimicry using vocab that includes mimicry n-grams (mimicry gap). + mimic_grams = _token_ngrams(mimic, 3) + stide_poisoned = StideModel(n=3, ngrams=set(window) | set(mimic_grams)) + mimic_mismatch, _ = stide_poisoned.score_window(mimic_grams) + + passed = ( + len(raw_tokens) >= 20 + and len(normal_tokens) >= 12 + and markov.ready + and float(m_score) > float(n_score) + and float(mimic_mismatch) < 0.05 + ) + return { + "pass": passed, + "socket": path, + "events": len(raw_tokens), + "normal_events": len(normal_tokens), + "ngrams_window": len(window), + "markov_ready": markov.ready, + "markov_normal": n_score, + "markov_mimicry": m_score, + "stide_mimicry_mismatch": round(float(mimic_mismatch), 4), + "stide_live_vocab": len(stide.ngrams), + } diff --git a/kernel_ai/ml/rules/sigma/privesc_euid.json b/kernel_ai/ml/rules/sigma/privesc_euid.json new file mode 100644 index 0000000..b426b11 --- /dev/null +++ b/kernel_ai/ml/rules/sigma/privesc_euid.json @@ -0,0 +1,16 @@ +{ + "id": "privesc_euid", + "title": "Effective root with non-root real uid", + "description": "Process has euid=0 while ruid is non-root — elevation / setuid abuse candidate.", + "when": { + "kind_in": ["privesc"] + }, + "attack": { + "family": "privilege_escalation", + "mitre": "T1548", + "name": "Abuse Elevation Control Mechanism", + "confidence": 0.92, + "cve": [], + "color": "#e0564e" + } +} diff --git a/kernel_ai/ml/rules/sigma/reverse_shell_lineage.json b/kernel_ai/ml/rules/sigma/reverse_shell_lineage.json new file mode 100644 index 0000000..6a7bc44 --- /dev/null +++ b/kernel_ai/ml/rules/sigma/reverse_shell_lineage.json @@ -0,0 +1,18 @@ +{ + "id": "reverse_shell_lineage", + "title": "Shell spawned network utility", + "description": "Interactive shell launching nc/socat/curl often indicates reverse shell or staging.", + "when": { + "source_in": ["stage5_process"], + "parent_in": ["bash", "sh", "dash", "zsh"], + "child_in": ["nc", "ncat", "netcat", "socat", "curl", "wget"] + }, + "attack": { + "family": "command_and_control", + "mitre": "T1071", + "name": "Application Layer Protocol", + "confidence": 0.88, + "cve": [], + "color": "#67c8e0" + } +} diff --git a/kernel_ai/ml/rules/sigma/shell_from_web.json b/kernel_ai/ml/rules/sigma/shell_from_web.json new file mode 100644 index 0000000..6c9f0c4 --- /dev/null +++ b/kernel_ai/ml/rules/sigma/shell_from_web.json @@ -0,0 +1,18 @@ +{ + "id": "shell_from_web", + "title": "Web server spawned a shell", + "description": "Parent looks like a web/runtime worker; child is a shell or scripting interpreter (classic webshell/RCE pattern).", + "when": { + "source_in": ["stage5_process"], + "parent_in": ["nginx", "apache2", "httpd", "php-fpm", "php", "uwsgi", "gunicorn", "caddy", "node"], + "child_in": ["bash", "sh", "dash", "zsh", "python", "python3", "perl", "nc", "ncat", "netcat"] + }, + "attack": { + "family": "execution", + "mitre": "T1059", + "name": "Command and Scripting Interpreter", + "confidence": 0.9, + "cve": [], + "color": "#e8a54b" + } +} diff --git a/kernel_ai/ml/sequence_deep/__init__.py b/kernel_ai/ml/sequence_deep/__init__.py new file mode 100644 index 0000000..6ad2b0b --- /dev/null +++ b/kernel_ai/ml/sequence_deep/__init__.py @@ -0,0 +1,9 @@ +"""Stage 8 — deep sequence models (HMM → LSTM/Transformer). Stub package. + +Real training/inference lands later; this module locks the API so Stage 4/6 +windows can plug in without reshaping the worker. See ``docs/ML_STAGE8.md``. +""" + +from kernel_ai.ml.sequence_deep.scorer import DeepSequenceScorer + +__all__ = ["DeepSequenceScorer"] diff --git a/kernel_ai/ml/sequence_deep/__main__.py b/kernel_ai/ml/sequence_deep/__main__.py new file mode 100644 index 0000000..023ba3a --- /dev/null +++ b/kernel_ai/ml/sequence_deep/__main__.py @@ -0,0 +1,65 @@ +"""CLI: ``python -m kernel_ai.ml.sequence_deep markov [--synthetic|--corpus PATH]``.""" + +from __future__ import annotations + +import argparse +import json +import logging +import sys + +from kernel_ai.ml.config import MLConfig + + +def main(argv: list[str] | None = None) -> int: + logging.basicConfig(level=logging.INFO, format="%(levelname)s %(name)s %(message)s") + parser = argparse.ArgumentParser(description="Stage 8 deep-sequence training") + parser.add_argument("backend", choices=("markov", "lstm", "transformer")) + parser.add_argument("--corpus", type=str, default=None, help="text corpus (one seq/line)") + parser.add_argument( + "--synthetic", + action="store_true", + help="train Markov on built-in normal sequences (local/CI)", + ) + parser.add_argument( + "--min-transitions", + type=int, + default=50, + help="refuse to write artifact below this many transitions", + ) + parser.add_argument("--no-ngrams", action="store_true", help="do not read ml_syscall_ngrams") + args = parser.parse_args(argv) + cfg = MLConfig() + + if args.backend != "markov": + from kernel_ai.ml.sequence_deep.lstm import train_lstm_stub + + try: + train_lstm_stub(cfg) + except SystemExit as exc: + if isinstance(exc.code, str): + print(exc.code, file=sys.stderr) + return 2 + raise + return 0 + + from kernel_ai.ml.sequence_deep.train_markov import train_markov + + try: + metrics = train_markov( + cfg, + corpus_path=args.corpus, + use_ngrams=not args.no_ngrams and not args.synthetic and not args.corpus, + use_synthetic=bool(args.synthetic), + min_transitions=args.min_transitions, + ) + except SystemExit as exc: + if isinstance(exc.code, str): + print(exc.code, file=sys.stderr) + return 2 + raise + print(json.dumps(metrics, indent=2, default=str)) + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/kernel_ai/ml/sequence_deep/encode.py b/kernel_ai/ml/sequence_deep/encode.py new file mode 100644 index 0000000..53b9453 --- /dev/null +++ b/kernel_ai/ml/sequence_deep/encode.py @@ -0,0 +1,66 @@ +"""Syscall / n-gram → integer id encoding (Stage 8 stub). + +Future: fit a vocab on normal windows (ml_syscall_ngrams / ADFA-LD), persist +beside the model artifact, pad/truncate online windows for Markov/LSTM. +""" + +from __future__ import annotations + +from dataclasses import dataclass, field + + +@dataclass +class SequenceEncoder: + """Bidirectional token ↔ id map. Untrained stub has an empty vocab.""" + + unk_token: str = "" + pad_token: str = "" + token_to_id: dict[str, int] = field(default_factory=dict) + id_to_token: dict[int, str] = field(default_factory=dict) + + def __post_init__(self) -> None: + if not self.token_to_id: + self.token_to_id = {self.pad_token: 0, self.unk_token: 1} + self.id_to_token = {0: self.pad_token, 1: self.unk_token} + + @property + def ready(self) -> bool: + # More than pad/unk means a real vocab was loaded/fitted. + return len(self.token_to_id) > 2 + + def fit(self, tokens: list[str]) -> None: + for tok in tokens: + if tok not in self.token_to_id: + idx = len(self.token_to_id) + self.token_to_id[tok] = idx + self.id_to_token[idx] = tok + + def encode(self, tokens: list[str], *, length: int | None = None) -> list[int]: + unk = self.token_to_id[self.unk_token] + ids = [self.token_to_id.get(t, unk) for t in tokens] + if length is None: + return ids + if len(ids) >= length: + return ids[-length:] + pad = self.token_to_id[self.pad_token] + return [pad] * (length - len(ids)) + ids + + def decode(self, ids: list[int]) -> list[str]: + return [self.id_to_token.get(i, self.unk_token) for i in ids] + + def state_dict(self) -> dict: + return { + "unk_token": self.unk_token, + "pad_token": self.pad_token, + "token_to_id": dict(self.token_to_id), + } + + @classmethod + def from_state(cls, state: dict) -> "SequenceEncoder": + enc = cls( + unk_token=state.get("unk_token", ""), + pad_token=state.get("pad_token", ""), + ) + enc.token_to_id = dict(state.get("token_to_id") or enc.token_to_id) + enc.id_to_token = {int(v): k for k, v in enc.token_to_id.items()} + return enc diff --git a/kernel_ai/ml/sequence_deep/lstm.py b/kernel_ai/ml/sequence_deep/lstm.py new file mode 100644 index 0000000..2943183 --- /dev/null +++ b/kernel_ai/ml/sequence_deep/lstm.py @@ -0,0 +1,47 @@ +"""LSTM / Transformer language-model scorer (Stage 8 stub). + +Heavy deps (``torch``) must stay lazy — never import at module top level. +Untrained stub always reports ``ready=False`` and refuses to score. +""" + +from __future__ import annotations + +from dataclasses import dataclass, field + + +@dataclass +class LstmScorer: + """Placeholder for a small causal LM over syscall tokens.""" + + backend: str = "lstm" # lstm | transformer + meta: dict = field(default_factory=dict) + _loaded: bool = False + + @property + def ready(self) -> bool: + return self._loaded + + def load(self, _path: str) -> None: + """Load a torch checkpoint. Stub: marks not ready, no torch import.""" + # Future: + # import torch + # self._model = torch.load(path, map_location="cpu") + # self._loaded = True + self._loaded = False + raise FileNotFoundError( + f"Stage 8 {self.backend} artifact not available (stub). " + "Train offline after Stage 6 data exists — docs/ML_STAGE8.md" + ) + + def score_window(self, _token_ids: list[int]) -> dict | None: + if not self.ready: + return None + return None + + +def train_lstm_stub(_cfg) -> dict: + raise SystemExit( + "Stage 8 LSTM/Transformer training not implemented yet. " + "Requires torch (lazy), normal-only corpora, and Stage 6 fidelity. " + "See docs/ML_STAGE8.md." + ) diff --git a/kernel_ai/ml/sequence_deep/markov.py b/kernel_ai/ml/sequence_deep/markov.py new file mode 100644 index 0000000..07d2225 --- /dev/null +++ b/kernel_ai/ml/sequence_deep/markov.py @@ -0,0 +1,88 @@ +"""Markov / HMM sequence scorer (Stage 8). + +Order-1 transition table over syscall / n-gram tokens. Offline training lives +in :mod:`kernel_ai.ml.sequence_deep.train_markov`. Pure Python (no hmmlearn yet). +""" + +from __future__ import annotations + +import math +from dataclasses import dataclass, field + + +@dataclass +class MarkovScorer: + """Order-1 transition table. Untrained → :meth:`score_window` returns None.""" + + order: int = 1 + meta: dict = field(default_factory=dict) + # counts[prev][nxt] = count + _counts: dict[str, dict[str, int]] = field(default_factory=dict) + _row_totals: dict[str, int] = field(default_factory=dict) + + @property + def ready(self) -> bool: + return bool(self._counts) + + def observe(self, tokens: list[str]) -> None: + if len(tokens) < 2: + return + for a, b in zip(tokens, tokens[1:]): + row = self._counts.setdefault(a, {}) + row[b] = row.get(b, 0) + 1 + self._row_totals[a] = self._row_totals.get(a, 0) + 1 + + def score_window(self, tokens: list[str]) -> dict | None: + """Return neg-avg-logprob style score, or None if untrained / too short.""" + if not self.ready or len(tokens) < 2: + return None + total = 0.0 + n = 0 + worst_i = 1 + worst_lp = 0.0 + for i in range(1, len(tokens)): + prev, nxt = tokens[i - 1], tokens[i] + row_total = self._row_totals.get(prev, 0) + # Laplace-ish smoothing so unseen transitions are finite but rare. + vocab = max(1, len(self._counts)) + cnt = (self._counts.get(prev) or {}).get(nxt, 0) + prob = (cnt + 1.0) / (row_total + vocab) if row_total else 1.0 / vocab + lp = math.log(max(prob, 1e-12)) + total += lp + n += 1 + if lp < worst_lp: + worst_lp = lp + worst_i = i + if n <= 0: + return None + avg_lp = total / n + return { + "model": "markov", + "neg_avg_logprob": round(-avg_lp, 4), + "avg_logprob": round(avg_lp, 4), + "worst_index": worst_i, + "worst_tokens": tokens[max(0, worst_i - 1) : worst_i + 1], + "window_len": len(tokens), + } + + def state_dict(self) -> dict: + return { + "order": self.order, + "meta": dict(self.meta), + "counts": self._counts, + "row_totals": self._row_totals, + } + + @classmethod + def from_state(cls, state: dict) -> "MarkovScorer": + m = cls(order=int(state.get("order") or 1), meta=dict(state.get("meta") or {})) + m._counts = {k: dict(v) for k, v in (state.get("counts") or {}).items()} + m._row_totals = {k: int(v) for k, v in (state.get("row_totals") or {}).items()} + return m + + +def train_markov_stub(cfg) -> dict: + """Backward-compatible name → real trainer (prefer ``train_markov``).""" + from kernel_ai.ml.sequence_deep.train_markov import train_markov + + return train_markov(cfg, use_synthetic=True) diff --git a/kernel_ai/ml/sequence_deep/scorer.py b/kernel_ai/ml/sequence_deep/scorer.py new file mode 100644 index 0000000..efff877 --- /dev/null +++ b/kernel_ai/ml/sequence_deep/scorer.py @@ -0,0 +1,118 @@ +"""Facade used by the ML worker for Stage 8 online scoring (stub-safe).""" + +from __future__ import annotations + +import logging +import os +from typing import Any + +from kernel_ai.ml.sequence_deep.encode import SequenceEncoder +from kernel_ai.ml.sequence_deep.markov import MarkovScorer + +logger = logging.getLogger("kernel_ai.ml.sequence_deep") + + +class DeepSequenceScorer: + """Hot-reloadable Stage 8 scorer. + + Prefer Markov artifact when present; LSTM path is wired but stubbed. + If nothing is trained, :meth:`score_tokens` returns None (no mutations). + """ + + def __init__(self, cfg: Any) -> None: + self.cfg = cfg + self.encoder = SequenceEncoder() + self.markov = MarkovScorer() + self.lstm = None + self._markov_mtime: float | None = None + self._lstm_mtime: float | None = None + self.maybe_reload() + + def maybe_reload(self) -> None: + path = getattr(self.cfg, "stage8_markov_path", "") or "" + if path: + try: + mtime = os.path.getmtime(path) + except OSError: + mtime = None + if mtime is not None and (self._markov_mtime is None or mtime > self._markov_mtime): + try: + import joblib + + state = joblib.load(path) + if isinstance(state, dict) and "encoder" in state: + self.encoder = SequenceEncoder.from_state(state["encoder"]) + self.markov = MarkovScorer.from_state(state.get("markov") or {}) + else: + self.markov = MarkovScorer.from_state(state if isinstance(state, dict) else {}) + self._markov_mtime = mtime + logger.info("loaded Stage 8 Markov artifact: %s", path) + except Exception as exc: # noqa: BLE001 + logger.warning("Stage 8 Markov load failed: %s", exc) + + lstm_path = getattr(self.cfg, "stage8_lstm_path", "") or "" + if lstm_path and os.path.exists(lstm_path): + try: + mtime = os.path.getmtime(lstm_path) + except OSError: + return + if self._lstm_mtime is not None and mtime <= self._lstm_mtime: + return + try: + from kernel_ai.ml.sequence_deep.lstm import LstmScorer + + scorer = LstmScorer(backend=getattr(self.cfg, "stage8_backend", "lstm")) + scorer.load(lstm_path) + self.lstm = scorer + self._lstm_mtime = mtime + except Exception as exc: # noqa: BLE001 + logger.info("Stage 8 LSTM idle (stub/untrained): %s", exc) + + @property + def ready(self) -> bool: + return self.markov.ready or (self.lstm is not None and self.lstm.ready) + + def score_tokens(self, tokens: list[str]) -> dict | None: + """Score an ordered token window (syscall names or n-gram keys).""" + self.maybe_reload() + if not tokens: + return None + if self.markov.ready: + return self.markov.score_window(tokens) + if self.lstm is not None and self.lstm.ready: + ids = self.encoder.encode(tokens, length=getattr(self.cfg, "stage8_window", 64)) + return self.lstm.score_window(ids) + return None + + def build_anomaly(self, score: dict, cfg: Any) -> dict: + """Map a Stage 8 score dict onto the shared mutation contract.""" + neg = float(score.get("neg_avg_logprob") or score.get("perplexity") or 0.0) + warn = float(getattr(cfg, "stage8_score_warn", 3.0)) + crit = float(getattr(cfg, "stage8_score_crit", 5.0)) + severity = "high" if neg >= crit else "medium" + worst = score.get("worst_tokens") or [] + why = " → ".join(str(t) for t in worst) if worst else score.get("model", "deep-seq") + return { + "source": "stage8_sequence", + "feature": "syscall_seq_deep", + "subsystem": "sched", + "type": "syscall_sequence_deep", + "severity": severity, + "score": round(neg, 4), + "value": float(score.get("window_len") or 0), + "baseline_mean": None, + "baseline_std": None, + "position": 0.18, + "message": ( + f"Deep sequence model ({score.get('model', '?')}): " + f"neg_avg_logprob={neg:.2f} (warn={warn}); unlikely transition near {why}" + ), + "meta": { + "stage": 8, + "model": score.get("model"), + "perplexity": score.get("perplexity"), + "neg_avg_logprob": neg, + "worst_index": score.get("worst_index"), + "worst_tokens": worst, + }, + } diff --git a/kernel_ai/ml/sequence_deep/train_markov.py b/kernel_ai/ml/sequence_deep/train_markov.py new file mode 100644 index 0000000..75cdc32 --- /dev/null +++ b/kernel_ai/ml/sequence_deep/train_markov.py @@ -0,0 +1,159 @@ +"""Offline Markov training for Stage 8. + +Corpus sources (first that yields enough transitions wins, unless forced): + 1. ``--corpus`` text file — one sequence per line (tokens separated by space or ``|``) + 2. Postgres ``ml_syscall_ngrams`` — expand n-gram keys by their counts + 3. ``--synthetic`` — built-in “normal” hostish sequences (local demo / CI) + +Writes ``cfg.stage8_markov_path`` joblib: ``{encoder, markov}`` for worker hot-reload. +""" + +from __future__ import annotations + +import logging +import os +from pathlib import Path + +import joblib + +from kernel_ai.ml.config import MLConfig +from kernel_ai.ml.sequence_deep.encode import SequenceEncoder +from kernel_ai.ml.sequence_deep.markov import MarkovScorer + +logger = logging.getLogger("kernel_ai.ml.sequence_deep.train") + +# Quiet, repetitive patterns a normal host tends to emit (demo prior only). +_SYNTHETIC_NORMAL = [ + ["read", "read", "write", "read", "close"], + ["futex", "futex", "poll", "futex"], + ["recvfrom", "recvfrom", "sendto", "recvfrom"], + ["epoll_wait", "epoll_wait", "read", "write", "epoll_wait"], + ["mmap", "munmap", "mmap", "munmap"], + ["openat", "read", "read", "close"], + ["openat", "fstat", "read", "close"], + ["clone", "futex", "futex", "exit"], + ["rt_sigaction", "rt_sigprocmask", "nanosleep"], + ["getpid", "gettid", "clock_gettime"], + ["stat", "openat", "read", "close", "stat"], + ["write", "write", "fdatasync"], +] + + +def _parse_line(line: str) -> list[str]: + line = line.strip() + if not line or line.startswith("#"): + return [] + if "|" in line and " " not in line.split("|")[0]: + return [t for t in line.split("|") if t] + return [t for t in line.split() if t] + + +def load_corpus_file(path: str | Path) -> list[list[str]]: + sequences: list[list[str]] = [] + with open(path, "r", encoding="utf-8", errors="ignore") as fh: + for line in fh: + seq = _parse_line(line) + if len(seq) >= 2: + sequences.append(seq) + return sequences + + +def load_corpus_ngrams(dsn: str, *, n: int, min_count: int = 1) -> list[list[str]]: + from kernel_ai.ml.store import fetch_ngram_counts + + counts = fetch_ngram_counts(dsn, n=n) + sequences: list[list[str]] = [] + for key, cnt in counts.items(): + if cnt < min_count: + continue + toks = [t for t in str(key).split("|") if t] + if len(toks) < 2: + continue + # Cap repeats so a hot n-gram cannot dominate the table entirely. + reps = min(int(cnt), 50) + for _ in range(reps): + sequences.append(toks) + return sequences + + +def load_corpus_synthetic(*, repeats: int = 40) -> list[list[str]]: + sequences: list[list[str]] = [] + for _ in range(max(1, repeats)): + sequences.extend(seq[:] for seq in _SYNTHETIC_NORMAL) + return sequences + + +def _count_transitions(sequences: list[list[str]]) -> int: + return sum(max(0, len(s) - 1) for s in sequences) + + +def train_markov( + cfg: MLConfig | None = None, + *, + corpus_path: str | None = None, + use_ngrams: bool = True, + use_synthetic: bool = False, + min_transitions: int = 50, +) -> dict: + """Fit Markov + encoder and persist artifact. Returns metrics dict.""" + cfg = cfg or MLConfig() + sequences: list[list[str]] = [] + source = "empty" + + if corpus_path: + sequences = load_corpus_file(corpus_path) + source = f"file:{corpus_path}" + elif use_synthetic: + sequences = load_corpus_synthetic() + source = "synthetic" + elif use_ngrams: + try: + sequences = load_corpus_ngrams(cfg.dsn, n=cfg.seq_n, min_count=cfg.seq_min_ngram_count) + source = "ml_syscall_ngrams" + except Exception as exc: # noqa: BLE001 + logger.warning("ngram corpus unavailable (%s) — falling back to synthetic", exc) + sequences = load_corpus_synthetic() + source = "synthetic_fallback" + + n_trans = _count_transitions(sequences) + if n_trans < min_transitions: + raise SystemExit( + f"Not enough transitions to train Markov (have {n_trans}, need >={min_transitions}, source={source}). " + "Provide --corpus, accumulate Stage 4/6 n-grams, or pass --synthetic." + ) + + encoder = SequenceEncoder() + markov = MarkovScorer(order=1, meta={"stage": 8, "source": source}) + for seq in sequences: + encoder.fit(seq) + markov.observe(seq) + + artifact = { + "encoder": encoder.state_dict(), + "markov": markov.state_dict(), + } + out_path = cfg.stage8_markov_path + os.makedirs(os.path.dirname(out_path) or ".", exist_ok=True) + joblib.dump(artifact, out_path) + + # Quick self-check: normal-ish window vs weird jump. + normal_score = markov.score_window(_SYNTHETIC_NORMAL[0]) + weird_score = markov.score_window(["openat", "execve", "connect", "dup2"]) + metrics = { + "source": source, + "n_sequences": len(sequences), + "n_transitions": n_trans, + "vocab": len(encoder.token_to_id), + "states": len(markov._counts), + "path": out_path, + "normal_neg_avg_logprob": (normal_score or {}).get("neg_avg_logprob"), + "weird_neg_avg_logprob": (weird_score or {}).get("neg_avg_logprob"), + } + logger.info( + "saved Stage 8 Markov → %s (source=%s transitions=%d vocab=%d)", + out_path, + source, + n_trans, + metrics["vocab"], + ) + return metrics diff --git a/tests/test_ml_collectors.py b/tests/test_ml_collectors.py new file mode 100644 index 0000000..a5a4a33 --- /dev/null +++ b/tests/test_ml_collectors.py @@ -0,0 +1,32 @@ +"""Tests for Stage 6 syscall event contract + n-gram stream ingest.""" + +from kernel_ai.ml.collectors.base import SyscallEvent, decode_events, encode_events +from kernel_ai.ml.sequence import NgramTracker + + +def test_encode_decode_roundtrip(): + events = [ + SyscallEvent(ts=1.0, pid=10, uid=0, comm="bash", syscall="clone"), + SyscallEvent(ts=1.1, pid=10, uid=0, comm="bash", syscall="execve"), + ] + payload = encode_events(events) + got = decode_events(payload) + assert len(got) == 2 + assert got[0].syscall == "clone" + assert got[1].pid == 10 + + +def test_ngram_tracker_update_stream_builds_trigrams(): + tracker = NgramTracker(n=3, window=50) + events = [ + SyscallEvent(ts=1.0, pid=7, uid=0, comm="x", syscall="clone"), + SyscallEvent(ts=1.1, pid=7, uid=0, comm="x", syscall="openat"), + SyscallEvent(ts=1.2, pid=7, uid=0, comm="x", syscall="execve"), + SyscallEvent(ts=1.3, pid=7, uid=0, comm="x", syscall="connect"), + ] + assert tracker.update_stream(events) == 4 + recent = tracker.recent() + assert "clone|openat|execve" in recent + assert "openat|execve|connect" in recent + pending = tracker.drain_pending() + assert pending["clone|openat|execve"] == 1 diff --git a/tests/test_ml_polygon.py b/tests/test_ml_polygon.py new file mode 100644 index 0000000..ed81a22 --- /dev/null +++ b/tests/test_ml_polygon.py @@ -0,0 +1,38 @@ +"""Local Stage 5/7/8 polygon tests.""" + +from kernel_ai.ml.collectors.stream_e2e import run_stream_e2e +from kernel_ai.ml.polygon import run_dry, run_mimicry + + +def test_polygon_dry_run_all_pass(): + results = run_dry( + [ + "reverse_shell", + "web_shell", + "miner_stub", + "scanner", + "privesc", + "lineage_shell", + ], + write_labels=False, + ) + failed = [r["scenario"] for r in results if not r["pass"]] + assert not failed, failed + + +def test_polygon_mimicry_stide_miss_markov_hit(): + result = run_mimicry() + assert result["stide"]["misses_mimicry"] is True + assert result["stide"]["mimicry_mismatch"] == 0.0 + assert result["markov"]["catches_mimicry"] is True + assert result["markov"]["mimicry_neg_avg_logprob"] > result["markov"]["normal_neg_avg_logprob"] + assert result["stide"]["catches_novel"] is True + assert result["pass"] is True + + +def test_stage6_stream_e2e_socket(): + result = run_stream_e2e(bursts=9, demo_every=0.02) + assert result["events"] >= 20, result + assert result["markov_ready"] is True + assert result["markov_mimicry"] > result["markov_normal"] + assert result["pass"] is True diff --git a/tests/test_ml_stage5.py b/tests/test_ml_stage5.py new file mode 100644 index 0000000..cce1850 --- /dev/null +++ b/tests/test_ml_stage5.py @@ -0,0 +1,74 @@ +"""Stage 5 process/lineage detector unit tests.""" + +from kernel_ai.ml.proc_baseline import LineageWhitelist, ProcBaselineDetector +from kernel_ai.ml.proc_features import ProcSample + + +def _sample(**kwargs) -> ProcSample: + base = dict( + pid=1000, + ppid=1, + comm="sleep", + parent_comm="bash", + ruid=1000, + euid=1000, + age_sec=5.0, + num_threads=1, + fd_count=4, + vm_rss_mb=2.0, + ) + base.update(kwargs) + s = ProcSample(**base) + s.features = s.score_vector() + return s + + +def test_lineage_whitelist_poison_guard(): + wl = LineageWhitelist(min_count=3) + assert wl.observe("bash", "nc") == 1 + assert wl.is_below_threshold("bash", "nc") + wl.observe("bash", "nc") + assert wl.is_below_threshold("bash", "nc") + wl.observe("bash", "nc") + assert not wl.is_below_threshold("bash", "nc") + + +def test_detector_emits_novel_lineage_once_per_pid(): + det = ProcBaselineDetector( + alpha=0.1, + warmup_samples=0, + z_warn=4.0, + z_crit=7.0, + lineage_min_count=3, + cooldown_sec=0.0, + max_emit_per_tick=8, + ) + s1 = _sample(pid=42, comm="nc", parent_comm="bash", age_sec=5.0) + out1 = det.score([s1], now=100.0) + assert any(a["source"] == "stage5_process" and "lineage:" in a["feature"] for a in out1) + # Same pid again: no second lineage observe/alert. + out2 = det.score([s1], now=101.0) + assert not any(a.get("meta", {}).get("kind") == "lineage" for a in out2) + + +def test_detector_privesc_rule(): + det = ProcBaselineDetector( + alpha=0.1, + warmup_samples=100, # keep EWMA quiet + z_warn=4.0, + z_crit=7.0, + lineage_min_count=1, + cooldown_sec=0.0, + max_emit_per_tick=8, + ) + # Pre-seed lineage so the edge is already normal. + det.lineage.load_counts([("bash", "sudo", 10)]) + s = _sample(pid=7, comm="sudo", parent_comm="bash", ruid=1000, euid=0, age_sec=30.0) + out = det.score([s], now=50.0) + assert any(a["type"] == "proc_anomaly:euid_root" for a in out) + + +def test_stage5_default_off(): + from kernel_ai.ml.config import MLConfig + + assert MLConfig().enable_stage5 is False diff --git a/tests/test_ml_stage7.py b/tests/test_ml_stage7.py new file mode 100644 index 0000000..b451009 --- /dev/null +++ b/tests/test_ml_stage7.py @@ -0,0 +1,66 @@ +"""Stage 7 attribution enricher tests.""" + +from kernel_ai.ml.attribution.attack_map import map_anomaly +from kernel_ai.ml.attribution.enrich import enrich_anomaly, enrich_anomalies +from kernel_ai.ml.attribution.sigma_engine import load_rules, match_anomaly + + +def test_sigma_rules_load(): + rules = load_rules() + ids = {r["id"] for r in rules} + assert "shell_from_web" in ids + assert "reverse_shell_lineage" in ids + + +def test_sigma_reverse_shell_rule(): + anom = { + "source": "stage5_process", + "feature": "lineage:bash->nc", + "type": "lineage:bash->nc", + "message": "Unusual process lineage", + "meta": {"kind": "lineage", "comm": "nc", "parent_comm": "bash"}, + } + hit = match_anomaly(anom) + assert hit is not None + assert hit["mitre"] == "T1071" + assert hit["source"] == "sigma" + + +def test_heuristic_pgfault_maps_to_impact(): + attack = map_anomaly( + { + "source": "stage1_baseline", + "feature": "pgfault_per_sec", + "type": "baseline_spike:pgfault_per_sec", + "message": "minor faults spike", + } + ) + assert attack is not None + assert attack["mitre"] == "T1499" + + +def test_enrich_writes_meta_and_message_prefix(): + anom = { + "source": "stage5_process", + "feature": "lineage:nginx->bash", + "type": "lineage:nginx->bash", + "message": "Unusual process lineage: nginx → bash", + "meta": {"kind": "lineage", "comm": "bash", "parent_comm": "nginx"}, + } + out = enrich_anomaly(anom, min_confidence=0.3) + assert out["attack"]["mitre"] == "T1059" + assert out["meta"]["attack"]["mitre"] == "T1059" + assert out["message"].startswith("[T1059]") + + +def test_enrich_batch_respects_min_confidence(): + weak = { + "source": "stage2_isoforest", + "feature": "vector", + "type": "isoforest:vector", + "message": "IsolationForest", + "meta": {}, + } + # heuristic confidence for bare isoforest is 0.30 → filtered at 0.35 + out = enrich_anomalies([weak], min_confidence=0.35) + assert "attack" not in out[0] diff --git a/tests/test_ml_stage8.py b/tests/test_ml_stage8.py new file mode 100644 index 0000000..b7decfd --- /dev/null +++ b/tests/test_ml_stage8.py @@ -0,0 +1,119 @@ +"""Stage 8 deep-sequence tests (Markov train + scorer).""" + +from types import SimpleNamespace + +import joblib + +from kernel_ai.ml.config import MLConfig +from kernel_ai.ml.sequence_deep.encode import SequenceEncoder +from kernel_ai.ml.sequence_deep.markov import MarkovScorer +from kernel_ai.ml.sequence_deep.scorer import DeepSequenceScorer +from kernel_ai.ml.sequence_deep.train_markov import train_markov + + +def test_stage8_default_off(): + assert MLConfig().enable_stage8 is False + + +def test_encoder_fit_encode_roundtrip(): + enc = SequenceEncoder() + assert not enc.ready + enc.fit(["read", "write", "close"]) + assert enc.ready + ids = enc.encode(["read", "write", "unknown"], length=4) + assert ids[0] == enc.token_to_id[""] + assert enc.decode(ids)[-2] == "write" + + +def test_markov_scores_after_observe(): + m = MarkovScorer() + assert m.score_window(["a", "b"]) is None + m.observe(["read", "write", "close", "read", "write", "close"]) + assert m.ready + normal = m.score_window(["read", "write", "close"]) + weird = m.score_window(["read", "execve", "connect"]) + assert normal is not None and weird is not None + assert weird["neg_avg_logprob"] >= normal["neg_avg_logprob"] + + +def test_train_markov_synthetic_and_reload(tmp_path): + out = tmp_path / "markov_latest.joblib" + cfg = MLConfig() + # dataclass is frozen — build a tiny namespace with required fields + ns = SimpleNamespace( + dsn=cfg.dsn, + seq_n=cfg.seq_n, + seq_min_ngram_count=cfg.seq_min_ngram_count, + stage8_markov_path=str(out), + ) + metrics = train_markov(ns, use_synthetic=True, use_ngrams=False, min_transitions=50) + assert out.is_file() + assert metrics["n_transitions"] >= 50 + assert metrics["weird_neg_avg_logprob"] >= metrics["normal_neg_avg_logprob"] + + scorer_cfg = SimpleNamespace( + stage8_markov_path=str(out), + stage8_lstm_path=str(tmp_path / "no.pt"), + stage8_backend="markov", + stage8_window=16, + stage8_score_warn=3.0, + stage8_score_crit=5.0, + ) + scorer = DeepSequenceScorer(scorer_cfg) + assert scorer.ready + score = scorer.score_tokens(["openat", "execve", "connect"]) + assert score is not None + assert score["model"] == "markov" + + +def test_deep_scorer_noop_without_artifact(): + cfg = SimpleNamespace( + stage8_markov_path="/tmp/kai-no-such-markov.joblib", + stage8_lstm_path="/tmp/kai-no-such-lstm.pt", + stage8_backend="markov", + stage8_window=16, + stage8_score_warn=3.0, + stage8_score_crit=5.0, + ) + scorer = DeepSequenceScorer(cfg) + assert not scorer.ready + assert scorer.score_tokens(["read", "write"]) is None + + +def test_build_anomaly_contract(): + cfg = SimpleNamespace( + stage8_markov_path="/tmp/kai-no-such-markov.joblib", + stage8_lstm_path="/tmp/no.pt", + stage8_backend="markov", + stage8_window=16, + stage8_score_warn=3.0, + stage8_score_crit=5.0, + ) + scorer = DeepSequenceScorer(cfg) + anom = scorer.build_anomaly( + { + "model": "markov", + "neg_avg_logprob": 4.2, + "worst_tokens": ["read", "execve"], + "window_len": 10, + }, + cfg, + ) + assert anom["source"] == "stage8_sequence" + assert anom["type"] == "syscall_sequence_deep" + assert anom["meta"]["stage"] == 8 + assert anom["severity"] == "medium" + + +def test_corpus_file_train(tmp_path): + corpus = tmp_path / "norm.txt" + corpus.write_text( + "\n".join(["read write close"] * 30 + ["futex futex poll"] * 20) + "\n", + encoding="utf-8", + ) + out = tmp_path / "m.joblib" + ns = SimpleNamespace(dsn="", seq_n=3, seq_min_ngram_count=1, stage8_markov_path=str(out)) + metrics = train_markov(ns, corpus_path=str(corpus), use_ngrams=False, min_transitions=40) + assert metrics["source"].startswith("file:") + blob = joblib.load(out) + assert "markov" in blob and "encoder" in blob