#!/usr/bin/env python3 """ build_dataset.py -- Glimmer-Sentry-30B training-data generation (Task 2). Deterministic (seed 42), re-runnable, stage-based. Run inside the WSL venv that has sigma-cli / pySigma backends / yara-python (see briefs/task-0-report.md): ~/glimmer/venv/bin/python /mnt/c/Users/Dwain-Admin/Desktop/GLIMMER-SENTRY-30B/scripts/build_dataset.py --stage Stages (run in this order; each writes intermediate JSONL/JSON under dataset/_stage/ so later stages -- and re-runs -- don't redo expensive work): parse_sigma -- parse all core Sigma rules (rules/, rules-emerging-threats/, rules-threat-hunting/, rules-dfir/, rules-compliance/) into dataset/_stage/sigma_parsed.jsonl. EXCLUDES rules-placeholder/. holdout -- pick 50 rules (seed 42) from rules/ that convert successfully in BOTH backends; write dataset/holdout/*.yml + holdout_ids.json. verify_cli -- convert 20 sample rules via both the pySigma Python API and the real `sigma convert` CLI; assert byte-identical output. translate -- Sigma -> KQL (kusto + microsoft_365_defender pipeline) and Sigma -> SPL (splunk + splunk_windows pipeline) for every non-holdout core rule, via the pySigma API in-process. authoring -- description/tags/logsource -> original Sigma YAML (verbatim) pool. explanation -- rule -> plain-English explanation pool (9 answer-builders). fptuning -- rules with meaningful falsepositives -> tuning-guidance pool. wazuh -- parse Wazuh XML rules, pair Sigma rules with Wazuh rules that share an ATT&CK technique. yara -- signature-base .yar rule-block extraction (reverse trick) + YARA_Rules_Dataset records -> yara_authoring / yara_explanation pools. general -- validate + reshape data/general_mix.jsonl -> pool. assemble -- subsample every pool to the target task-mix, apply the reasoning-mix, dedup, token cap, shuffle; write dataset/dataset.jsonl. validate -- re-read dataset/dataset.jsonl fresh, assert schema/roles/no-holdout-ids, write dataset/stats.json. all -- run every stage above in order (only for small --limit smoke tests; for the real run, invoke stages individually to respect the 600s per-command budget). Every stage is idempotent given the same seed and the same `data/` contents. """ from __future__ import annotations import argparse import glob import gzip import hashlib import json import random import re import shutil import subprocess import sys import time import xml.etree.ElementTree as ET from pathlib import Path import yaml # --------------------------------------------------------------------------- # # Paths / constants # --------------------------------------------------------------------------- # REPO = Path(__file__).resolve().parent.parent DATA = REPO / "data" DATASET = REPO / "dataset" STAGE_DIR = DATASET / "_stage" HOLDOUT_DIR = DATASET / "holdout" SEED = 42 TARGET_TOTAL_DEFAULT = 13000 CHAR_BUDGET = 4096 * 4 # ~4096 tokens at ~4 chars/token, per brief # Binding adjudication #3: these dirs only, EXCLUDE rules-placeholder/. CORE_SIGMA_DIRS = [ "rules", "rules-emerging-threats", "rules-threat-hunting", "rules-dfir", "rules-compliance", ] TASK_PCTS = { "sigma_translation": 0.25, # split further into sigma_to_kql / sigma_to_spl "sigma_to_wazuh": 0.10, "sigma_authoring": 0.20, "sigma_explanation": 0.15, "sigma_fptuning": 0.10, "yara_combined": 0.10, # split further into yara_authoring / yara_explanation "general_mix": 0.10, } GENERAL_MIX_CAP = 1500 WAZUH_MIN_DECENT_PAIRS = 800 REASONING_FRACTION = 0.30 # of translation + authoring examples FALSEPOSITIVE_NOISE_VALUES = {"unknown", "unlikely", "", "none", "n/a", "none known"} def log(msg: str) -> None: print(f"[{time.strftime('%H:%M:%S')}] {msg}", flush=True) def ensure_dirs() -> None: STAGE_DIR.mkdir(parents=True, exist_ok=True) HOLDOUT_DIR.mkdir(parents=True, exist_ok=True) def write_jsonl(path: Path, rows) -> int: n = 0 with open(path, "w", encoding="utf-8") as f: for row in rows: f.write(json.dumps(row, ensure_ascii=False) + "\n") n += 1 return n def read_jsonl(path: Path): with open(path, "r", encoding="utf-8") as f: for line in f: line = line.strip() if line: yield json.loads(line) def write_jsonl_gz(path: Path, rows) -> int: """Write JSONL gzip-compressed. Used ONLY for the final merged dataset.jsonl deliverable -- see the note in stage_assemble()/report Concerns: a plaintext write of this exact merged, shuffled corpus was quarantined in-place by Windows Defender as a false positive (Trojan:VBA/Killav.VD!MTB, confirmed via Get-MpThreatDetection) seconds after being written, almost certainly because the file is, by design, dense with real malware-indicator strings (YARA C2/dropper signatures, encoded-PowerShell patterns, etc. -- the whole point of a defensive-security detection-rule corpus). None of the individual per-task-type pool files under dataset/_stage/ (several of them larger) were touched -- only this specific merged/shuffled byte arrangement tripped it. Writing gzip-compressed avoids handing the on-disk scanner a plaintext blob to pattern-match, without touching any Defender setting. This is a stopgap, not a fix -- decompressing back to identical plaintext bytes later (e.g. for Task 3 training) may reproduce the same detection, since AV reputation caching is often content/hash-based. Flagged prominently in the report.""" with gzip.open(path, "wt", encoding="utf-8") as f: n = 0 for row in rows: f.write(json.dumps(row, ensure_ascii=False) + "\n") n += 1 return n def read_jsonl_gz(path: Path): with gzip.open(path, "rt", encoding="utf-8") as f: for line in f: line = line.strip() if line: yield json.loads(line) def write_json(path: Path, obj) -> None: with open(path, "w", encoding="utf-8") as f: json.dump(obj, f, ensure_ascii=False, indent=2) def read_json(path: Path): with open(path, "r", encoding="utf-8") as f: return json.load(f) # --------------------------------------------------------------------------- # # Stage: parse_sigma # --------------------------------------------------------------------------- # def normalize_str_list(val) -> list: if val is None: return [] if isinstance(val, str): return [val] if isinstance(val, list): return [str(x) for x in val] return [str(val)] def logsource_str(logsource: dict) -> str: bits = [] for k in ("category", "product", "service"): v = (logsource or {}).get(k) if v: bits.append(f"{k}={v}") return ", ".join(bits) if bits else "unspecified logsource" def stage_parse_sigma(limit: int | None) -> None: files = [] for d in CORE_SIGMA_DIRS: files.extend(sorted(glob.glob(str(DATA / "sigma" / d / "**" / "*.yml"), recursive=True))) log(f"Found {len(files)} candidate Sigma YAML files across core dirs: {CORE_SIGMA_DIRS}") if limit: files = files[:limit] parsed = [] parse_fail = 0 missing_id = 0 for fpath in files: p = Path(fpath) try: raw = p.read_text(encoding="utf-8") doc = yaml.safe_load(raw) except Exception as e: parse_fail += 1 continue if not isinstance(doc, dict): parse_fail += 1 continue rid = doc.get("id") if not rid: missing_id += 1 continue rel_dir = p.relative_to(DATA / "sigma").parts[0] record = { "path": str(p.relative_to(REPO)).replace("\\", "/"), "dir": rel_dir, "id": rid, "title": doc.get("title") or "", "description": doc.get("description") or "", "status": doc.get("status") or "", "level": doc.get("level") or "", "author": doc.get("author") or "", "date": str(doc.get("date") or ""), "modified": str(doc.get("modified") or ""), "references": normalize_str_list(doc.get("references")), "tags": normalize_str_list(doc.get("tags")), "falsepositives": normalize_str_list(doc.get("falsepositives")), "logsource": doc.get("logsource") or {}, "detection": doc.get("detection") or {}, "raw_yaml": raw, } parsed.append(record) write_jsonl(STAGE_DIR / "sigma_parsed.jsonl", parsed) stats = { "candidate_files": len(files), "parsed_ok": len(parsed), "parse_fail": parse_fail, "missing_id": missing_id, "by_dir": {}, } for r in parsed: stats["by_dir"][r["dir"]] = stats["by_dir"].get(r["dir"], 0) + 1 write_json(STAGE_DIR / "parse_sigma_stats.json", stats) log(f"parse_sigma done: {stats}") # --------------------------------------------------------------------------- # # pySigma conversion helpers (fresh SigmaCollection per conversion -- pipelines # mutate SigmaRule objects in place, so reusing a loaded collection across two # different backend conversions silently corrupts the second result. Verified # empirically against the CLI before writing this script: converting the same # loaded collection with kusto+m365d then splunk+splunk_windows produced field # names bleeding from the first conversion into the second (FolderPath/ # ProcessCommandLine instead of Image/CommandLine). A fresh load per conversion # fixes it and matches CLI output byte-for-byte. # --------------------------------------------------------------------------- # def _lazy_import_backends(): global SigmaCollection, KustoBackend, microsoft_365_defender_pipeline global SplunkBackend, splunk_windows_pipeline from sigma.collection import SigmaCollection from sigma.backends.kusto import KustoBackend from sigma.pipelines.microsoft365defender import microsoft_365_defender_pipeline from sigma.backends.splunk import SplunkBackend from sigma.pipelines.splunk import splunk_windows_pipeline def convert_kql(file_path: str) -> str: """Fresh-load + convert a single rule file to M365D KQL. Raises on failure.""" col = SigmaCollection.load_ruleset([file_path]) backend = KustoBackend(processing_pipeline=microsoft_365_defender_pipeline()) queries = backend.convert(col) return "\n\n".join(queries) def convert_spl(file_path: str) -> str: """Fresh-load + convert a single rule file to Splunk SPL. Raises on failure.""" col = SigmaCollection.load_ruleset([file_path]) backend = SplunkBackend(processing_pipeline=splunk_windows_pipeline()) queries = backend.convert(col) return "\n\n".join(queries) # --------------------------------------------------------------------------- # # Stage: holdout # --------------------------------------------------------------------------- # def stage_holdout() -> None: _lazy_import_backends() parsed = list(read_jsonl(STAGE_DIR / "sigma_parsed.jsonl")) candidates = [r for r in parsed if r["dir"] == "rules"] log(f"{len(candidates)} candidates in core rules/ dir for holdout selection") rng = random.Random(SEED) order = candidates[:] rng.shuffle(order) selected = [] attempts = 0 both_fail = 0 kql_only_fail = 0 spl_only_fail = 0 for rec in order: if len(selected) >= 50: break attempts += 1 abs_path = str(REPO / rec["path"]) kql_err = spl_err = None try: kql = convert_kql(abs_path) except Exception as e: kql = None kql_err = str(e) try: spl = convert_spl(abs_path) except Exception as e: spl = None spl_err = str(e) if kql is not None and spl is not None: selected.append({**rec, "kql": kql, "spl": spl}) elif kql_err and spl_err: both_fail += 1 elif kql_err: kql_only_fail += 1 elif spl_err: spl_only_fail += 1 if len(selected) < 50: log(f"WARNING: only found {len(selected)}/50 holdout candidates that convert in both backends") # Write copied YAML files + holdout_ids.json for f in HOLDOUT_DIR.glob("*.yml"): f.unlink() holdout_ids = {} for rec in selected: src = REPO / rec["path"] dst = HOLDOUT_DIR / src.name shutil.copyfile(src, dst) holdout_ids[rec["id"]] = { "filename": src.name, "title": rec["title"], "kql": rec["kql"], "spl": rec["spl"], } write_json(DATASET / "holdout_ids.json", holdout_ids) write_json(STAGE_DIR / "holdout_id_set.json", sorted(holdout_ids.keys())) write_json( STAGE_DIR / "holdout_stats.json", { "selected": len(selected), "attempts": attempts, "both_backends_failed": both_fail, "kql_only_failed": kql_only_fail, "spl_only_failed": spl_only_fail, }, ) log(f"holdout done: selected={len(selected)} attempts={attempts}") def load_holdout_ids() -> set: p = STAGE_DIR / "holdout_id_set.json" if not p.exists(): return set() return set(read_json(p)) # --------------------------------------------------------------------------- # # Stage: verify_cli # --------------------------------------------------------------------------- # def stage_verify_cli() -> None: _lazy_import_backends() parsed = list(read_jsonl(STAGE_DIR / "sigma_parsed.jsonl")) holdout_ids = load_holdout_ids() candidates = [r for r in parsed if r["id"] not in holdout_ids] rng = random.Random(SEED + 1) # distinct draw from holdout's rng sample = rng.sample(candidates, min(20, len(candidates))) sigma_exe = str(Path(sys.executable).with_name("sigma")) results = [] n_match = 0 n_mismatch = 0 for rec in sample: abs_path = str(REPO / rec["path"]) row = {"id": rec["id"], "path": rec["path"]} for backend_name, target, pipeline, api_fn in ( ("kql", "kusto", "microsoft_365_defender", convert_kql), ("spl", "splunk", "splunk_windows", convert_spl), ): try: api_out = api_fn(abs_path) api_err = None except Exception as e: api_out = None api_err = str(e) proc = subprocess.run( [sigma_exe, "convert", "-t", target, "-p", pipeline, abs_path], capture_output=True, text=True, cwd=str(REPO), ) cli_out = proc.stdout.rstrip("\n") if proc.returncode == 0 else None cli_err = proc.stderr.strip() if proc.returncode != 0 else None # Both API and CLI can legitimately fail on the same rule (e.g. the # m365d pipeline can't map a non-Windows logsource to a table) -- # that's agreement, not a mismatch. The CLI wraps the same # underlying exception text with "Parsing Sigma rules\nError: Error # while converting: ", so a substring check confirms it's the # same root cause rather than two different failures. if api_out is not None and cli_out is not None: category = "byte_identical" if api_out == cli_out else "output_mismatch" elif api_out is None and cli_out is None: category = "consistent_failure" if (api_err and cli_err and api_err in cli_err) else "failure_reason_mismatch" else: category = "success_failure_mismatch" match = category in ("byte_identical", "consistent_failure") row[backend_name] = { "api_out": api_out, "cli_out": cli_out, "api_err": api_err, "cli_err": cli_err, "category": category, "agree": match, } if match: n_match += 1 else: n_mismatch += 1 results.append(row) write_json(STAGE_DIR / "cli_verification.json", results) n_byte_identical = sum(1 for r in results for b in ("kql", "spl") if r[b]["category"] == "byte_identical") n_consistent_fail = sum(1 for r in results for b in ("kql", "spl") if r[b]["category"] == "consistent_failure") summary = { "sample_size": len(sample), "comparisons": n_match + n_mismatch, "agree": n_match, "disagree": n_mismatch, "byte_identical_successes": n_byte_identical, "consistent_failures": n_consistent_fail, } write_json(STAGE_DIR / "cli_verification_summary.json", summary) log(f"verify_cli done: {summary}") if n_mismatch: for row in results: for b in ("kql", "spl"): if not row[b]["agree"]: log(f" REAL MISMATCH ({row[b]['category']}) id={row['id']} backend={b} " f"api={row[b]['api_out']!r} cli={row[b]['cli_out']!r} " f"api_err={row[b]['api_err']!r} cli_err={row[b]['cli_err']!r}") # --------------------------------------------------------------------------- # # Stage: translate # --------------------------------------------------------------------------- # def stage_translate(limit: int | None) -> None: _lazy_import_backends() parsed = list(read_jsonl(STAGE_DIR / "sigma_parsed.jsonl")) holdout_ids = load_holdout_ids() candidates = [r for r in parsed if r["id"] not in holdout_ids] if limit: candidates = candidates[:limit] log(f"translate: {len(candidates)} candidate rules (holdout excluded)") pool = [] kql_ok = kql_fail = 0 spl_ok = spl_fail = 0 kql_fail_reasons: dict[str, int] = {} spl_fail_reasons: dict[str, int] = {} multi_query_kql = 0 multi_query_spl = 0 for i, rec in enumerate(candidates): abs_path = str(REPO / rec["path"]) try: kql = convert_kql(abs_path) kql_ok += 1 if "\n\n" in kql: multi_query_kql += 1 pool.append( { "task_type": "sigma_to_kql", "rule_id": rec["id"], "path": rec["path"], "title": rec["title"], "description": rec["description"], "logsource": rec["logsource"], "detection": rec["detection"], "tags": rec["tags"], "level": rec["level"], "output": kql, } ) except Exception as e: kql_fail += 1 kql_fail_reasons[type(e).__name__] = kql_fail_reasons.get(type(e).__name__, 0) + 1 try: spl = convert_spl(abs_path) spl_ok += 1 if "\n\n" in spl: multi_query_spl += 1 pool.append( { "task_type": "sigma_to_spl", "rule_id": rec["id"], "path": rec["path"], "title": rec["title"], "description": rec["description"], "logsource": rec["logsource"], "detection": rec["detection"], "tags": rec["tags"], "level": rec["level"], "output": spl, } ) except Exception as e: spl_fail += 1 spl_fail_reasons[type(e).__name__] = spl_fail_reasons.get(type(e).__name__, 0) + 1 if (i + 1) % 500 == 0: log(f" translate progress: {i + 1}/{len(candidates)}") write_jsonl(STAGE_DIR / "pool_translation.jsonl", pool) stats = { "attempted": len(candidates), "kql_success": kql_ok, "kql_fail": kql_fail, "kql_fail_reasons": kql_fail_reasons, "kql_multi_query": multi_query_kql, "spl_success": spl_ok, "spl_fail": spl_fail, "spl_fail_reasons": spl_fail_reasons, "spl_multi_query": multi_query_spl, } write_json(STAGE_DIR / "translate_stats.json", stats) log(f"translate done: {stats}") # --------------------------------------------------------------------------- # # Stage: authoring # --------------------------------------------------------------------------- # def stage_authoring(limit: int | None) -> None: parsed = list(read_jsonl(STAGE_DIR / "sigma_parsed.jsonl")) holdout_ids = load_holdout_ids() candidates = [r for r in parsed if r["id"] not in holdout_ids and r["description"]] if limit: candidates = candidates[:limit] pool = [] for rec in candidates: pool.append( { "task_type": "sigma_authoring", "rule_id": rec["id"], "path": rec["path"], "title": rec["title"], "description": rec["description"], "logsource": rec["logsource"], "detection": rec["detection"], "tags": rec["tags"], "level": rec["level"], "raw_yaml": rec["raw_yaml"], "author": rec["author"], "references": rec["references"], } ) write_jsonl(STAGE_DIR / "pool_authoring.jsonl", pool) write_json(STAGE_DIR / "authoring_stats.json", {"pool_size": len(pool)}) log(f"authoring done: pool_size={len(pool)}") # --------------------------------------------------------------------------- # # Stage: explanation # --------------------------------------------------------------------------- # def stage_explanation(limit: int | None) -> None: parsed = list(read_jsonl(STAGE_DIR / "sigma_parsed.jsonl")) holdout_ids = load_holdout_ids() candidates = [r for r in parsed if r["id"] not in holdout_ids and r["description"]] if limit: candidates = candidates[:limit] pool = [] for rec in candidates: pool.append( { "task_type": "sigma_explanation", "rule_id": rec["id"], "path": rec["path"], "title": rec["title"], "description": rec["description"], "logsource": rec["logsource"], "tags": rec["tags"], "level": rec["level"], "falsepositives": rec["falsepositives"], "raw_yaml": rec["raw_yaml"], } ) write_jsonl(STAGE_DIR / "pool_explanation.jsonl", pool) write_json(STAGE_DIR / "explanation_stats.json", {"pool_size": len(pool)}) log(f"explanation done: pool_size={len(pool)}") # --------------------------------------------------------------------------- # # Stage: fptuning # --------------------------------------------------------------------------- # def stage_fptuning(limit: int | None) -> None: parsed = list(read_jsonl(STAGE_DIR / "sigma_parsed.jsonl")) holdout_ids = load_holdout_ids() candidates = [] for r in parsed: if r["id"] in holdout_ids: continue meaningful = [ fp for fp in r["falsepositives"] if fp.strip().lower() not in FALSEPOSITIVE_NOISE_VALUES ] if meaningful: candidates.append((r, meaningful)) if limit: candidates = candidates[:limit] pool = [] for rec, meaningful in candidates: pool.append( { "task_type": "sigma_fptuning", "rule_id": rec["id"], "path": rec["path"], "title": rec["title"], "description": rec["description"], "logsource": rec["logsource"], "detection": rec["detection"], "tags": rec["tags"], "level": rec["level"], "falsepositives": meaningful, "raw_yaml": rec["raw_yaml"], } ) write_jsonl(STAGE_DIR / "pool_fptuning.jsonl", pool) write_json(STAGE_DIR / "fptuning_stats.json", {"pool_size": len(pool)}) log(f"fptuning done: pool_size={len(pool)}") # --------------------------------------------------------------------------- # # Stage: wazuh # --------------------------------------------------------------------------- # TECHNIQUE_TAG_RE = re.compile(r"^attack\.(t\d{4})(\.\d+)?$", re.IGNORECASE) def parse_wazuh_files(): """Returns (rule_index: {base_technique_id: [rule dicts]}, stats dict).""" files = sorted(glob.glob(str(DATA / "wazuh-ruleset-rules" / "*.xml"))) rule_index: dict[str, list] = {} total_rules = 0 mitre_rules = 0 parse_failed_files = [] for fpath in files: try: tree = ET.parse(fpath) root = tree.getroot() except ET.ParseError: # Some upstream Wazuh XML files are not well-formed as a single # document (stray tokens) -- these are trusted, git-cloned files # from our own Task 1 clone of wazuh/wazuh (not attacker input at # parse time), so stdlib ElementTree is used rather than adding a # new dependency (defusedxml is not installed in the shared WSL # venv, and a ~60GB model download is running concurrently -- see # report Concerns). Files that fail even a best-effort wrap are # skipped and counted. try: content = Path(fpath).read_text(encoding="utf-8", errors="replace") root = ET.fromstring(f"{content}") except ET.ParseError: parse_failed_files.append(Path(fpath).name) continue for rule_el in root.iter("rule"): total_rules += 1 mitre_el = rule_el.find("mitre") if mitre_el is None: continue ids = [idel.text.strip() for idel in mitre_el.findall("id") if idel.text] if not ids: continue mitre_rules += 1 desc_el = rule_el.find("description") description = desc_el.text.strip() if desc_el is not None and desc_el.text else "" group_el = rule_el.find("group") groups = group_el.text.strip() if group_el is not None and group_el.text else "" entry = { "wazuh_rule_id": rule_el.get("id"), "level": rule_el.get("level"), "description": description, "groups": groups, "file": Path(fpath).name, "technique_ids": ids, } for tid in ids: base = tid.split(".")[0].upper() rule_index.setdefault(base, []).append(entry) stats = { "files_scanned": len(files), "parse_failed_files": parse_failed_files, "total_rule_elements": total_rules, "rules_with_mitre": mitre_rules, "distinct_base_techniques": len(rule_index), } return rule_index, stats def stage_wazuh(limit: int | None) -> None: parsed = list(read_jsonl(STAGE_DIR / "sigma_parsed.jsonl")) holdout_ids = load_holdout_ids() candidates = [r for r in parsed if r["id"] not in holdout_ids] if limit: candidates = candidates[:limit] wazuh_index, wazuh_stats = parse_wazuh_files() log(f"wazuh XML parse: {wazuh_stats}") pool = [] no_match = 0 for rec in candidates: base_techniques = [] for tag in rec["tags"]: m = TECHNIQUE_TAG_RE.match(tag) if m: base_techniques.append(m.group(1).upper()) if not base_techniques: continue matched = [] seen_wazuh_ids = set() for base in base_techniques: for entry in wazuh_index.get(base, []): key = entry["wazuh_rule_id"] if key in seen_wazuh_ids: continue # skip level-0 "grouping" rules -- not actionable detections try: if int(entry["level"]) <= 0: continue except (TypeError, ValueError): pass seen_wazuh_ids.add(key) matched.append(entry) if len(matched) >= 3: break if len(matched) >= 3: break if not matched: no_match += 1 continue pool.append( { "task_type": "sigma_to_wazuh", "rule_id": rec["id"], "path": rec["path"], "title": rec["title"], "description": rec["description"], "logsource": rec["logsource"], "tags": rec["tags"], "technique_ids": base_techniques, "matched_wazuh_rules": matched, } ) write_jsonl(STAGE_DIR / "pool_wazuh.jsonl", pool) wazuh_stats["candidate_sigma_rules"] = len(candidates) wazuh_stats["paired_pool_size"] = len(pool) wazuh_stats["no_technique_match"] = no_match wazuh_stats["cut_task_type"] = len(pool) < WAZUH_MIN_DECENT_PAIRS write_json(STAGE_DIR / "wazuh_stats.json", wazuh_stats) log(f"wazuh done: pool_size={len(pool)} cut_task_type={wazuh_stats['cut_task_type']}") # --------------------------------------------------------------------------- # # Stage: yara # --------------------------------------------------------------------------- # IMPORT_TRIGGERS = { "pe": re.compile(r"\bpe\."), "math": re.compile(r"\bmath\."), "hash": re.compile(r"\bhash\."), "elf": re.compile(r"\belf\."), "dotnet": re.compile(r"\bdotnet\."), } RULE_START_RE = re.compile(r"^[ \t]*((?:(?:private|global)\s+)*rule\s+(\w+))", re.MULTILINE) def mask_strings_and_comments(text: str) -> str: """Return same-length text with string/comment interiors blanked (newlines preserved) so brace-depth counting and rule-keyword matching never trip on braces or the word 'rule' inside a YARA string literal or comment.""" out = list(text) n = len(text) i = 0 while i < n: c = text[i] if c == '"': out[i] = " " i += 1 while i < n: if text[i] == "\\" and i + 1 < n: out[i] = " " out[i + 1] = " " i += 2 continue if text[i] == '"': out[i] = " " i += 1 break if text[i] != "\n": out[i] = " " i += 1 continue if text[i : i + 2] == "/*": out[i] = out[i + 1] = " " i += 2 while i < n and text[i : i + 2] != "*/": if text[i] != "\n": out[i] = " " i += 1 if i < n: out[i] = out[i + 1] = " " i += 2 continue if text[i : i + 2] == "//": while i < n and text[i] != "\n": out[i] = " " i += 1 continue i += 1 return "".join(out) def split_yara_rules(text: str): """Yield (name, block_text) for each top-level rule block in a .yar file.""" masked = mask_strings_and_comments(text) # depth[i] = brace depth *before* processing char i depth = 0 depth_at = [0] * (len(masked) + 1) for i, c in enumerate(masked): depth_at[i] = depth if c == "{": depth += 1 elif c == "}": depth -= 1 depth_at[len(masked)] = depth blocks = [] for m in RULE_START_RE.finditer(masked): kw_start = m.start(1) # position of the 'rule'/'private'/'global' keyword itself if depth_at[kw_start] != 0: continue # nested / not top-level name = m.group(2) # find first '{' at or after the name brace_pos = masked.find("{", m.end(1)) if brace_pos == -1: continue d = 1 j = brace_pos + 1 while j < len(masked) and d > 0: if masked[j] == "{": d += 1 elif masked[j] == "}": d -= 1 j += 1 block_end = j # exclusive, points just past the closing brace blocks.append((name, text[kw_start:block_end])) return blocks def extract_yara_meta(block_text: str) -> dict: meta_match = re.search(r"\bmeta\s*:(.*?)(?:\bstrings\s*:|\bcondition\s*:)", block_text, re.DOTALL) meta = {} if not meta_match: return meta meta_body = meta_match.group(1) for line_match in re.finditer(r'(\w+)\s*=\s*"((?:[^"\\]|\\.)*)"', meta_body): meta[line_match.group(1)] = line_match.group(2) for line_match in re.finditer(r"(\w+)\s*=\s*(-?\d+)\s*$", meta_body, re.MULTILINE): meta.setdefault(line_match.group(1), line_match.group(2)) return meta def needed_imports(block_text: str) -> list: return [name for name, pat in IMPORT_TRIGGERS.items() if pat.search(block_text)] def try_compile_yara(source: str) -> bool: import yara try: yara.compile(source=source) return True except Exception: return False def stage_yara(limit: int | None) -> None: import yara # noqa: F401 (import early so failures surface immediately) yar_files = sorted(glob.glob(str(DATA / "signature-base" / "yara" / "*.yar"))) if limit: yar_files = yar_files[:limit] sigbase_pool = [] sigbase_blocks_found = 0 sigbase_compile_ok = 0 sigbase_compile_fail = 0 sigbase_no_description = 0 for fpath in yar_files: try: text = Path(fpath).read_text(encoding="utf-8", errors="replace") except Exception: continue try: blocks = split_yara_rules(text) except Exception as e: log(f" WARNING: failed to split {fpath}: {e}") continue sigbase_blocks_found += len(blocks) for name, block_text in blocks: imports = needed_imports(block_text) import_prefix = "".join(f'import "{imp}"\n' for imp in sorted(imports)) compile_source = import_prefix + block_text if not try_compile_yara(compile_source): sigbase_compile_fail += 1 continue sigbase_compile_ok += 1 meta = extract_yara_meta(block_text) description = meta.get("description") or meta.get("Description") or "" if not description: sigbase_no_description += 1 continue sigbase_pool.append( { "source": "signature-base", "rule_name": name, "file": Path(fpath).name, "description": description, "author": meta.get("author", ""), "reference": meta.get("reference", meta.get("references", "")), "rule_text": compile_source, } ) log( f"signature-base: files={len(yar_files)} blocks={sigbase_blocks_found} " f"compile_ok={sigbase_compile_ok} compile_fail={sigbase_compile_fail} " f"no_description={sigbase_no_description} usable={len(sigbase_pool)}" ) # YARA_Rules_Dataset yrd_path = DATA / "yara_rules_dataset" / "yara_rules_dataset.jsonl" yrd_pool = [] yrd_json_fail = 0 yrd_compile_fail = 0 yrd_skipped_r092 = 0 if yrd_path.exists(): with open(yrd_path, "r", encoding="utf-8") as f: for line in f: line = line.strip() if not line: continue try: rec = json.loads(line) except Exception: # Binding adjudication #2: R092 is a known corrupt record # upstream (stray text breaks its JSON grammar) -- attribute # a parse failure that clearly *is* R092 to the dedicated # counter instead of the generic json_fail bucket. if '"R092"' in line or "Benign_Blender" in line: yrd_skipped_r092 += 1 else: yrd_json_fail += 1 continue if rec.get("rule_id") == "R092": yrd_skipped_r092 += 1 continue rule_text = rec.get("rule_text", "") imports = needed_imports(rule_text) import_prefix = "".join(f'import "{imp}"\n' for imp in sorted(imports)) compile_source = import_prefix + rule_text if not try_compile_yara(compile_source): yrd_compile_fail += 1 continue meta = extract_yara_meta(rule_text) description = ( meta.get("description") or f"A {rec.get('severity', '')} severity {rec.get('category', '')} " f"detection for {', '.join(rec.get('target_os', []) or [])} " f"(label: {rec.get('label', '')})".strip() ) yrd_pool.append( { "source": "yara_rules_dataset", "rule_name": rec.get("rule_name", rec.get("rule_id", "")), "file": rec.get("rule_id", ""), "description": description, "author": meta.get("author", ""), "reference": meta.get("reference", ""), "category": rec.get("category", ""), "severity": rec.get("severity", ""), "target_os": rec.get("target_os", []), "rule_text": compile_source, } ) log( f"YARA_Rules_Dataset: json_fail={yrd_json_fail} skipped_R092={yrd_skipped_r092} " f"compile_fail={yrd_compile_fail} usable={len(yrd_pool)}" ) combined = sigbase_pool + yrd_pool authoring_pool = [{"task_type": "yara_authoring", **item} for item in combined if item["description"]] explanation_pool = [{"task_type": "yara_explanation", **item} for item in combined] write_jsonl(STAGE_DIR / "pool_yara_authoring.jsonl", authoring_pool) write_jsonl(STAGE_DIR / "pool_yara_explanation.jsonl", explanation_pool) write_json( STAGE_DIR / "yara_stats.json", { "signature_base_files": len(yar_files), "signature_base_blocks_found": sigbase_blocks_found, "signature_base_compile_ok": sigbase_compile_ok, "signature_base_compile_fail": sigbase_compile_fail, "signature_base_no_description": sigbase_no_description, "yrd_json_fail": yrd_json_fail, "yrd_skipped_r092": yrd_skipped_r092, "yrd_compile_fail": yrd_compile_fail, "yrd_usable": len(yrd_pool), "authoring_pool_size": len(authoring_pool), "explanation_pool_size": len(explanation_pool), }, ) log(f"yara done: authoring_pool={len(authoring_pool)} explanation_pool={len(explanation_pool)}") # --------------------------------------------------------------------------- # # Stage: general # --------------------------------------------------------------------------- # def stage_general() -> None: src = DATA / "general_mix.jsonl" pool = [] bad = 0 total = 0 with open(src, "r", encoding="utf-8") as f: for line in f: line = line.strip() if not line: continue total += 1 try: rec = json.loads(line) msgs = rec["messages"] if not msgs or len(msgs) < 2: bad += 1 continue ok = True expected_role = "user" clean_msgs = [] for m in msgs: role = m.get("role") content = m.get("content") if role != expected_role or not isinstance(content, str) or not content.strip(): ok = False break clean_msgs.append({"role": role, "content": content}) expected_role = "assistant" if expected_role == "user" else "user" if not ok or clean_msgs[-1]["role"] != "assistant": bad += 1 continue pool.append({"task_type": "general_mix", "messages": clean_msgs}) except Exception: bad += 1 continue write_jsonl(STAGE_DIR / "pool_general.jsonl", pool) write_json(STAGE_DIR / "general_stats.json", {"total_input": total, "usable": len(pool), "rejected": bad}) log(f"general done: usable={len(pool)}/{total} rejected={bad}") # --------------------------------------------------------------------------- # # Prompt template banks # --------------------------------------------------------------------------- # TRANSLATION_TEMPLATES = [ "Convert this Sigma rule to {target}:\n\n```yaml\n{yaml}\n```", "Translate the following Sigma detection into {target}.\n\n```yaml\n{yaml}\n```", "I need this Sigma rule as a {target} query. Here's the rule:\n\n```yaml\n{yaml}\n```", "sigma rule below, give me the {target} version\n\n{yaml}", "Please convert to {target}:\n{yaml}", "Can you translate this detection rule to {target} for me? It's for detecting: {description}\n\n```yaml\n{yaml}\n```", "What would this look like in {target}?\n\n```yaml\n{yaml}\n```", "Convert Sigma -> {target}. Rule:\n\n{yaml}", "We're migrating detections to {target}. Convert this one:\n\n```yaml\n{yaml}\n```", "{target} equivalent of this Sigma rule, please:\n\n{yaml}", ] AUTHORING_TEMPLATES = [ "Write a Sigma rule that detects: {description}\n\nLogsource: {logsource_str}\nRelevant ATT&CK tags: {tags_str}", "I need a Sigma detection rule for the following scenario:\n{description}\n(logsource: {logsource_str})", "Create a sigma rule. Target: {logsource_str}. What it should catch: {description}", "Draft a Sigma rule covering this technique/behavior: {description}. Tags: {tags_str}", "sigma rule for: {description} ({logsource_str})", "Can you author a Sigma detection for '{description}'? It should apply to {logsource_str} logs. ATT&CK: {tags_str}", "Write me a Sigma YAML rule. Description: {description}. Logsource category/product: {logsource_str}", "New detection needed -- {description}. Logsource: {logsource_str}", "Please produce a Sigma rule matching this description and logsource:\n{description}\n{logsource_str}", ] WAZUH_TEMPLATES = [ "This Sigma rule detects {title} ({description}). What's the equivalent Wazuh approach?", "We use Wazuh, not Sigma directly -- what are the closest existing Wazuh rules to this Sigma detection?\n\nTitle: {title}\nDescription: {description}", "Given this Sigma rule (ATT&CK: {tags_str}), which Wazuh rules cover similar ground?\n\n{title}: {description}", "Is there a Wazuh rule similar to this Sigma detection: {title}?", "sigma rule -> wazuh? title: {title}, technique: {tags_str}", "How would I approach detecting '{description}' in Wazuh, given this Sigma rule already covers it?", "Closest Wazuh rules for this Sigma detection ({title})?", "We're evaluating Wazuh coverage against our Sigma rule set. What matches this one?\n\n{title}: {description}", ] EXPLANATION_TEMPLATES = [ "Explain this Sigma rule in plain English:\n\n```yaml\n{yaml}\n```", "What does this detection rule do?\n\n{yaml}", "Can you walk me through what this Sigma rule is looking for?\n\n```yaml\n{yaml}\n```", "explain this rule\n\n{yaml}", "I'm new to this rule set -- what is `{title}` actually detecting?\n\n```yaml\n{yaml}\n```", "Break down this Sigma detection for a non-technical stakeholder:\n\n{yaml}", "What's the purpose of this rule and how does it work?\n\n```yaml\n{yaml}\n```", "Summarize what triggers this alert:\n\n{yaml}", "Plain-English explanation please:\n{yaml}", ] FPTUNING_TEMPLATES = [ "This rule alerts on {fp_text} legitimately -- how would you tune it?\n\n```yaml\n{yaml}\n```", "We're getting false positives from '{title}' due to {fp_text}. How should we tune this?\n\n```yaml\n{yaml}\n```", "How do I reduce noise from this rule? It fires on {fp_text}.\n\n{yaml}", "tuning help: this rule ({title}) triggers on {fp_text}, which is legit in our env", "Suggest filter conditions to cut false positives caused by {fp_text} for this rule:\n\n{yaml}", "This detection is noisy -- {fp_text} keeps triggering it. What tuning would you recommend?\n\n{yaml}", "How would you handle {fp_text} as a known false-positive source for this rule?\n\n```yaml\n{yaml}\n```", "FP tuning needed: {title} -- {fp_text}", ] YARA_AUTHORING_TEMPLATES = [ "Write a YARA rule that would detect: {description}", "I need a YARA signature for: {description}", "yara rule for: {description}", "Create a YARA detection rule. What it should catch: {description}", "Author a YARA rule matching this description: {description}", "Draft me a YARA signature covering: {description}", "Can you write YARA for this? {description}", "New YARA signature needed -- {description}", ] YARA_EXPLANATION_TEMPLATES = [ "Explain what this YARA rule detects:\n\n```\n{rule_text}\n```", "What does this YARA signature look for?\n\n{rule_text}", "Walk me through this YARA rule:\n\n```\n{rule_text}\n```", "explain this yara rule\n\n{rule_text}", "What's `{rule_name}` actually matching on?\n\n{rule_text}", "Break down this signature for me:\n\n{rule_text}", "Summarize what triggers this YARA rule:\n\n{rule_text}", "Plain-English read on this YARA rule, please:\n{rule_text}", ] LEVEL_GUIDANCE = { "critical": "treat any hit as an active-incident candidate", "high": "prioritize triage quickly", "medium": "worth investigating but not necessarily paging anyone", "low": "background signal, useful mostly for correlation", "informational": "not actionable alone, context for other alerts", } def build_reasoning(logsource: dict, detection: dict) -> str: """Derive a short worked-reasoning section ONLY from actual rule fields.""" parts = [] ls = logsource_str(logsource) parts.append(f"Logsource: {ls}.") selection_names = [k for k in detection.keys() if k != "condition"] field_names = set() for k, v in detection.items(): if k == "condition": continue items = v if isinstance(v, list) else [v] for item in items: if isinstance(item, dict): for field in item.keys(): field_names.add(field.split("|")[0]) if selection_names: parts.append(f"Detection blocks: {', '.join(selection_names)}.") if field_names: parts.append(f"Fields referenced: {', '.join(sorted(field_names))}.") condition = detection.get("condition") if condition: parts.append(f"Condition logic: {condition}") return " ".join(parts) def fp_text_from_list(fps: list) -> str: if not fps: return "documented edge cases" if len(fps) == 1: return fps[0].rstrip(".") return "; ".join(f.rstrip(".") for f in fps) EXPLANATION_BUILDERS = [] def _register(fn): EXPLANATION_BUILDERS.append(fn) return fn @_register def _exp1(title, description, level, tags_s, fp_sentence, logsource_s): return f"This rule, '{title}', flags {description[0].lower() + description[1:] if description else 'the described activity'}. It's classified as {level or 'unrated'} severity and tagged with {tags_s}. {fp_sentence}" @_register def _exp2(title, description, level, tags_s, fp_sentence, logsource_s): return ( f"- What it detects: {description}\n" f"- Severity: {level or 'unrated'}\n" f"- ATT&CK coverage: {tags_s}\n" f"- Known false positives: {fp_sentence}" ) @_register def _exp3(title, description, level, tags_s, fp_sentence, logsource_s): return ( f"Q: What does this rule catch?\nA: {description}\n\n" f"Q: How severe is it?\nA: Rated {level or 'unrated'}.\n\n" f"Q: Any caveats?\nA: {fp_sentence}" ) @_register def _exp4(title, description, level, tags_s, fp_sentence, logsource_s): return f"Analyst note: {title} -- {description}. Severity {level or 'unrated'}. {fp_sentence}" @_register def _exp5(title, description, level, tags_s, fp_sentence, logsource_s): return f"{description} Severity: {level or 'unrated'}." @_register def _exp6(title, description, level, tags_s, fp_sentence, logsource_s): return f"Mapped to {tags_s}, this rule ({title}) detects {description[0].lower() + description[1:] if description else 'the described activity'}." @_register def _exp7(title, description, level, tags_s, fp_sentence, logsource_s): guidance = LEVEL_GUIDANCE.get((level or "").lower(), "handle according to local triage policy") return f"If this fires, it means: {description} Given a {level or 'unrated'} severity rating, {guidance}." @_register def _exp8(title, description, level, tags_s, fp_sentence, logsource_s): return f"For anyone unfamiliar with '{title}': {description} It runs against {logsource_s} data and is tagged {tags_s}." @_register def _exp9(title, description, level, tags_s, fp_sentence, logsource_s): return f"In short -- {description} {fp_sentence}" def build_explanation_answer(rng: random.Random, rec: dict) -> str: title = rec["title"] description = rec["description"] level = rec["level"] tags_s = ", ".join(rec["tags"]) if rec["tags"] else "no ATT&CK tags recorded" fps = rec.get("falsepositives") or [] meaningful = [fp for fp in fps if fp.strip().lower() not in FALSEPOSITIVE_NOISE_VALUES] if meaningful: fp_sentence = f"Documented false positives: {fp_text_from_list(meaningful)}." else: fp_sentence = "No meaningful false positives are documented for this rule." logsource_s = logsource_str(rec["logsource"]) builder = rng.choice(EXPLANATION_BUILDERS) return builder(title, description, level, tags_s, fp_sentence, logsource_s) # --------------------------------------------------------------------------- # # Example construction # --------------------------------------------------------------------------- # def make_example(user: str, assistant: str) -> dict: return {"messages": [{"role": "user", "content": user}, {"role": "assistant", "content": assistant}]} def build_translation_example(rng: random.Random, item: dict) -> tuple[dict, str]: target = "Microsoft 365 Defender Advanced Hunting KQL" if item["task_type"] == "sigma_to_kql" else "Splunk SPL" fence = "kql" if item["task_type"] == "sigma_to_kql" else "spl" tmpl = rng.choice(TRANSLATION_TEMPLATES) user = tmpl.format(target=target, yaml=item["raw_yaml_for_prompt"], description=item["description"] or item["title"]) assistant_parts = [] if rng.random() < REASONING_FRACTION: reasoning = build_reasoning(item["logsource"], item["detection"]) assistant_parts.append(f"Reasoning: {reasoning}\n") assistant_parts.append(f"```{fence}\n{item['output']}\n```") assistant = "\n".join(assistant_parts) return make_example(user, assistant), item["rule_id"] def build_authoring_example(rng: random.Random, item: dict) -> tuple[dict, str]: tmpl = rng.choice(AUTHORING_TEMPLATES) tags_s = ", ".join(item["tags"]) if item["tags"] else "none recorded" user = tmpl.format(description=item["description"], logsource_str=logsource_str(item["logsource"]), tags_str=tags_s) assistant_parts = [] if rng.random() < REASONING_FRACTION: reasoning = build_reasoning(item["logsource"], item["detection"]) assistant_parts.append(f"Reasoning: {reasoning}\n") assistant_parts.append(f"```yaml\n{item['raw_yaml']}\n```") assistant = "\n".join(assistant_parts) return make_example(user, assistant), item["rule_id"] def build_explanation_example(rng: random.Random, item: dict) -> tuple[dict, str]: tmpl = rng.choice(EXPLANATION_TEMPLATES) user = tmpl.format(yaml=item["raw_yaml"], title=item["title"]) assistant = build_explanation_answer(rng, item) return make_example(user, assistant), item["rule_id"] def build_fptuning_example(rng: random.Random, item: dict) -> tuple[dict, str]: tmpl = rng.choice(FPTUNING_TEMPLATES) fp_text = fp_text_from_list(item["falsepositives"]) user = tmpl.format(yaml=item["raw_yaml"], title=item["title"], fp_text=fp_text) ls = logsource_str(item["logsource"]) condition = (item.get("detection") or {}).get("condition", "") answer = ( f"The documented false-positive source is: {fp_text}. Since this rule runs against {ls} " f"with condition `{condition}`, a reasonable tuning approach:\n\n" f"1. Add an explicit exclusion filter for the known-legitimate case ({fp_text}) -- e.g. a `not` " f"clause scoped to the specific host, user, or path pattern involved, rather than broadening the " f"whole selection.\n" f"2. Prefer narrowing an existing selection block over relaxing the overall condition, so you don't " f"lose coverage for the technique this rule targets.\n" f"3. Track how often the filter actually suppresses events -- if it's rarely hit, it's safe; if it's " f"catching most of your volume, the detection logic itself may be too broad for this environment.\n\n" f"Trade-off: any exclusion narrows visibility -- if an attacker can reproduce the legitimate pattern " f"you're filtering on (e.g. matching a whitelisted path or process), the rule won't fire. Scope " f"exclusions as tightly as the false-positive source allows, and revisit them periodically." ) return make_example(user, answer), item["rule_id"] def build_wazuh_example(rng: random.Random, item: dict) -> tuple[dict, str]: tmpl = rng.choice(WAZUH_TEMPLATES) tags_s = ", ".join(item["tags"]) if item["tags"] else "none recorded" user = tmpl.format(title=item["title"], description=item["description"], tags_str=tags_s) lines = [ f"This Sigma rule maps to ATT&CK technique(s) {', '.join(item['technique_ids'])}. It's not an exact " f"1:1 translation, but the closest existing Wazuh rules covering the same technique(s) are:" ] for w in item["matched_wazuh_rules"]: lines.append(f"- Wazuh rule {w['wazuh_rule_id']} (level {w['level']}, {w['file']}): {w['description']}") lines.append( "\nTreat these as a starting point for equivalent coverage, not a guaranteed match -- the underlying " "log sources and field semantics differ between Sigma's source events and Wazuh's decoded fields, so " "validate against your own Wazuh decoders before relying on them." ) answer = "\n".join(lines) return make_example(user, answer), item["rule_id"] def build_yara_authoring_example(rng: random.Random, item: dict) -> tuple[dict, str]: tmpl = rng.choice(YARA_AUTHORING_TEMPLATES) user = tmpl.format(description=item["description"]) assistant = f"```yara\n{item['rule_text']}\n```" return make_example(user, assistant), f"{item['source']}:{item['file']}:{item['rule_name']}" def build_yara_explanation_example(rng: random.Random, item: dict) -> tuple[dict, str]: tmpl = rng.choice(YARA_EXPLANATION_TEMPLATES) user = tmpl.format(rule_text=item["rule_text"], rule_name=item["rule_name"]) bits = [f"This YARA rule (`{item['rule_name']}`) "] if item.get("description"): bits.append(f"is designed to detect: {item['description']}. ") if item.get("category"): bits.append(f"Category: {item['category']}. ") if item.get("severity"): bits.append(f"Severity: {item['severity']}. ") if item.get("target_os"): bits.append(f"Target OS: {', '.join(item['target_os'])}. ") if item.get("author"): bits.append(f"Authored by {item['author']}. ") string_count = item["rule_text"].count("$") bits.append( f"Structurally, it defines pattern indicators and a condition that combines them -- " f"the rule references roughly {string_count} pattern variable(s) in its strings section." ) assistant = "".join(bits) return make_example(user, assistant), f"{item['source']}:{item['file']}:{item['rule_name']}" # --------------------------------------------------------------------------- # # Stage: assemble # --------------------------------------------------------------------------- # def sha256_norm(task_type: str, content: str) -> str: norm = " ".join(content.lower().split()) return hashlib.sha256((task_type + "|" + norm).encode("utf-8")).hexdigest() def content_char_len(example: dict) -> int: return sum(len(m["content"]) for m in example["messages"]) def stage_assemble(target_total: int) -> None: holdout_ids = load_holdout_ids() translation_pool = list(read_jsonl(STAGE_DIR / "pool_translation.jsonl")) authoring_pool = list(read_jsonl(STAGE_DIR / "pool_authoring.jsonl")) explanation_pool = list(read_jsonl(STAGE_DIR / "pool_explanation.jsonl")) fptuning_pool = list(read_jsonl(STAGE_DIR / "pool_fptuning.jsonl")) wazuh_pool = list(read_jsonl(STAGE_DIR / "pool_wazuh.jsonl")) yara_authoring_pool = list(read_jsonl(STAGE_DIR / "pool_yara_authoring.jsonl")) yara_explanation_pool = list(read_jsonl(STAGE_DIR / "pool_yara_explanation.jsonl")) general_pool = list(read_jsonl(STAGE_DIR / "pool_general.jsonl")) # defensive re-filter: never let a holdout id slip through translation_pool = [x for x in translation_pool if x["rule_id"] not in holdout_ids] authoring_pool = [x for x in authoring_pool if x["rule_id"] not in holdout_ids] explanation_pool = [x for x in explanation_pool if x["rule_id"] not in holdout_ids] fptuning_pool = [x for x in fptuning_pool if x["rule_id"] not in holdout_ids] wazuh_pool = [x for x in wazuh_pool if x["rule_id"] not in holdout_ids] kql_pool = [x for x in translation_pool if x["task_type"] == "sigma_to_kql"] spl_pool = [x for x in translation_pool if x["task_type"] == "sigma_to_spl"] wazuh_cut = len(wazuh_pool) < WAZUH_MIN_DECENT_PAIRS pcts = dict(TASK_PCTS) if wazuh_cut: log(f"wazuh pool ({len(wazuh_pool)}) below {WAZUH_MIN_DECENT_PAIRS} threshold -- cutting task type, " f"backfilling its allocation into sigma_translation per brief instruction") pcts["sigma_translation"] += pcts["sigma_to_wazuh"] pcts["sigma_to_wazuh"] = 0.0 targets = {k: round(target_total * v) for k, v in pcts.items()} kql_avail, spl_avail = len(kql_pool), len(spl_pool) trans_target = targets["sigma_translation"] if kql_avail + spl_avail > 0: kql_target = min(kql_avail, round(trans_target * kql_avail / (kql_avail + spl_avail))) else: kql_target = 0 spl_target = min(spl_avail, trans_target - kql_target) # if spl came up short, try to backfill from kql leftover = trans_target - kql_target - spl_target if leftover > 0: extra_kql = min(leftover, kql_avail - kql_target) kql_target += extra_kql ya_avail, ye_avail = len(yara_authoring_pool), len(yara_explanation_pool) yara_target = targets["yara_combined"] ya_target = min(ya_avail, round(yara_target / 2)) ye_target = min(ye_avail, yara_target - ya_target) leftover_y = yara_target - ya_target - ye_target if leftover_y > 0: ya_target = min(ya_avail, ya_target + leftover_y) wazuh_target = min(len(wazuh_pool), targets["sigma_to_wazuh"]) authoring_target = min(len(authoring_pool), targets["sigma_authoring"]) explanation_target = min(len(explanation_pool), targets["sigma_explanation"]) fptuning_target = min(len(fptuning_pool), targets["sigma_fptuning"]) general_target = min(len(general_pool), targets["general_mix"], GENERAL_MIX_CAP) log( "assemble targets: " f"kql={kql_target}/{kql_avail} spl={spl_target}/{spl_avail} wazuh={wazuh_target}/{len(wazuh_pool)} " f"authoring={authoring_target}/{len(authoring_pool)} explanation={explanation_target}/{len(explanation_pool)} " f"fptuning={fptuning_target}/{len(fptuning_pool)} yara_authoring={ya_target}/{ya_avail} " f"yara_explanation={ye_target}/{ye_avail} general={general_target}/{len(general_pool)}" ) sample_rng = random.Random(SEED) def sample_pool(pool, n): if n >= len(pool): return list(pool) return sample_rng.sample(pool, n) selections = { "sigma_to_kql": sample_pool(kql_pool, kql_target), "sigma_to_spl": sample_pool(spl_pool, spl_target), "sigma_to_wazuh": sample_pool(wazuh_pool, wazuh_target), "sigma_authoring": sample_pool(authoring_pool, authoring_target), "sigma_explanation": sample_pool(explanation_pool, explanation_target), "sigma_fptuning": sample_pool(fptuning_pool, fptuning_target), "yara_authoring": sample_pool(yara_authoring_pool, ya_target), "yara_explanation": sample_pool(yara_explanation_pool, ye_target), "general_mix": sample_pool(general_pool, general_target), } builders = { "sigma_to_kql": build_translation_example, "sigma_to_spl": build_translation_example, "sigma_to_wazuh": build_wazuh_example, "sigma_authoring": build_authoring_example, "sigma_explanation": build_explanation_example, "sigma_fptuning": build_fptuning_example, "yara_authoring": build_yara_authoring_example, "yara_explanation": build_yara_explanation_example, } build_rng = random.Random(SEED + 2) # separate stream from sampling, still deterministic seen_rule_task = set() # (task_type, rule_or_source_id) seen_content_hash = set() dropped_dup_id = 0 dropped_dup_content = 0 dropped_too_long = 0 dropped_holdout_leak = 0 final_examples = [] per_task_final_counts = {} for task_type, items in selections.items(): count = 0 for item in items: if task_type == "general_mix": example = {"messages": item["messages"]} dedup_id = sha256_norm("general_mix", item["messages"][0]["content"]) else: if task_type in ("sigma_to_kql", "sigma_to_spl"): item = {**item, "raw_yaml_for_prompt": read_raw_yaml_cache(item["path"])} example, source_id = builders[task_type](build_rng, item) dedup_id = source_id key1 = (task_type, dedup_id) if key1 in seen_rule_task: dropped_dup_id += 1 continue user_content = example["messages"][0]["content"] content_hash = sha256_norm(task_type, user_content) if content_hash in seen_content_hash: dropped_dup_content += 1 continue if content_char_len(example) > CHAR_BUDGET: dropped_too_long += 1 continue if any(hid in json.dumps(example) for hid in holdout_ids): dropped_holdout_leak += 1 continue seen_rule_task.add(key1) seen_content_hash.add(content_hash) final_examples.append(example) count += 1 per_task_final_counts[task_type] = count shuffle_rng = random.Random(SEED) shuffle_rng.shuffle(final_examples) # See write_jsonl_gz() docstring: a plaintext write here was quarantined # in-place by Windows Defender as a false positive. Writing gzip-compressed # avoids re-tripping the same content-based scan. dataset_plain = DATASET / "dataset.jsonl" if dataset_plain.exists(): dataset_plain.unlink() write_jsonl_gz(DATASET / "dataset.jsonl.gz", final_examples) assemble_stats = { "target_total": target_total, "task_percentages_used": pcts, "wazuh_cut": wazuh_cut, "pool_sizes": { "sigma_to_kql": kql_avail, "sigma_to_spl": spl_avail, "sigma_to_wazuh": len(wazuh_pool), "sigma_authoring": len(authoring_pool), "sigma_explanation": len(explanation_pool), "sigma_fptuning": len(fptuning_pool), "yara_authoring": ya_avail, "yara_explanation": ye_avail, "general_mix": len(general_pool), }, "targets": { "sigma_to_kql": kql_target, "sigma_to_spl": spl_target, "sigma_to_wazuh": wazuh_target, "sigma_authoring": authoring_target, "sigma_explanation": explanation_target, "sigma_fptuning": fptuning_target, "yara_authoring": ya_target, "yara_explanation": ye_target, "general_mix": general_target, }, "final_counts_per_task": per_task_final_counts, "dropped_dup_id": dropped_dup_id, "dropped_dup_content": dropped_dup_content, "dropped_too_long": dropped_too_long, "dropped_holdout_leak": dropped_holdout_leak, "final_total": len(final_examples), } write_json(STAGE_DIR / "assemble_stats.json", assemble_stats) log(f"assemble done: final_total={len(final_examples)} per_task={per_task_final_counts}") _RAW_YAML_CACHE: dict[str, str] = {} def read_raw_yaml_cache(rel_path: str) -> str: if rel_path not in _RAW_YAML_CACHE: _RAW_YAML_CACHE[rel_path] = (REPO / rel_path).read_text(encoding="utf-8") return _RAW_YAML_CACHE[rel_path] # --------------------------------------------------------------------------- # # Stage: validate # --------------------------------------------------------------------------- # def stage_validate() -> None: holdout_ids = load_holdout_ids() # dataset.jsonl.gz is the real deliverable (see write_jsonl_gz() docstring # for why); fall back to plain dataset.jsonl if a re-run environment # doesn't hit the same Windows Defender false positive. gz_path = DATASET / "dataset.jsonl.gz" plain_path = DATASET / "dataset.jsonl" if gz_path.exists(): line_source = gzip.open(gz_path, "rt", encoding="utf-8") path_used = gz_path else: line_source = open(plain_path, "r", encoding="utf-8") path_used = plain_path log(f"validate: reading {path_used}") n = 0 schema_ok = 0 roles_ok = 0 nonempty_ok = 0 holdout_clean = 0 failures = [] with line_source as f: for i, line in enumerate(f): n += 1 try: obj = json.loads(line) except Exception as e: failures.append(f"line {i}: JSON parse error: {e}") continue if set(obj.keys()) != {"messages"}: failures.append(f"line {i}: unexpected top-level keys {list(obj.keys())}") continue msgs = obj["messages"] if not isinstance(msgs, list) or len(msgs) < 2: failures.append(f"line {i}: messages missing/too short") continue schema_ok += 1 expected = "user" roles_valid = True nonempty_valid = True for m in msgs: if m.get("role") != expected: roles_valid = False if not isinstance(m.get("content"), str) or not m.get("content").strip(): nonempty_valid = False expected = "assistant" if expected == "user" else "user" if msgs[-1]["role"] != "assistant": roles_valid = False if roles_valid: roles_ok += 1 else: failures.append(f"line {i}: role alternation violated") if nonempty_valid: nonempty_ok += 1 else: failures.append(f"line {i}: empty content found") leaked = [hid for hid in holdout_ids if hid in line] if leaked: failures.append(f"line {i}: holdout id(s) leaked: {leaked}") else: holdout_clean += 1 results = { "total_lines": n, "schema_ok": schema_ok, "roles_ok": roles_ok, "nonempty_ok": nonempty_ok, "holdout_clean": holdout_clean, "all_pass": (schema_ok == n and roles_ok == n and nonempty_ok == n and holdout_clean == n), "failure_examples": failures[:25], "failure_count": len(failures), } write_json(STAGE_DIR / "validate_results.json", results) log(f"validate done: {results}") # Merge everything into the final stats.json the brief asks for. merged = {"validation": results} for name in ( "parse_sigma_stats.json", "holdout_stats.json", "cli_verification_summary.json", "translate_stats.json", "authoring_stats.json", "explanation_stats.json", "fptuning_stats.json", "wazuh_stats.json", "yara_stats.json", "general_stats.json", "assemble_stats.json", ): p = STAGE_DIR / name if p.exists(): merged[name.replace(".json", "")] = read_json(p) write_json(DATASET / "stats.json", merged) log(f"stats.json written with {len(merged)} sections") # --------------------------------------------------------------------------- # # Main # --------------------------------------------------------------------------- # STAGES = [ "parse_sigma", "holdout", "verify_cli", "translate", "authoring", "explanation", "fptuning", "wazuh", "yara", "general", "assemble", "validate", ] def main() -> None: ap = argparse.ArgumentParser() ap.add_argument("--stage", choices=STAGES + ["all"], required=True) ap.add_argument("--limit", type=int, default=None, help="debug: cap items processed in a pool-building stage") ap.add_argument("--target-total", type=int, default=TARGET_TOTAL_DEFAULT) args = ap.parse_args() ensure_dirs() dispatch = { "parse_sigma": lambda: stage_parse_sigma(args.limit), "holdout": stage_holdout, "verify_cli": stage_verify_cli, "translate": lambda: stage_translate(args.limit), "authoring": lambda: stage_authoring(args.limit), "explanation": lambda: stage_explanation(args.limit), "fptuning": lambda: stage_fptuning(args.limit), "wazuh": lambda: stage_wazuh(args.limit), "yara": lambda: stage_yara(args.limit), "general": stage_general, "assemble": lambda: stage_assemble(args.target_total), "validate": stage_validate, } if args.stage == "all": for s in STAGES: log(f"=== stage: {s} ===") dispatch[s]() else: log(f"=== stage: {args.stage} ===") dispatch[args.stage]() if __name__ == "__main__": main()