diff --git a/salt/elastalert/files/modules/so/securityonion-es.py b/salt/elastalert/files/modules/so/securityonion-es.py index 4ed91d15f..67ac699f8 100644 --- a/salt/elastalert/files/modules/so/securityonion-es.py +++ b/salt/elastalert/files/modules/so/securityonion-es.py @@ -9,6 +9,7 @@ from datetime import datetime from time import gmtime, strftime import hashlib +import ipaddress import re import requests,json from elastalert.alerts import Alerter, DateTimeEncoder @@ -35,6 +36,9 @@ class SecurityOnionESAlerter(Alerter): 'value_median': 'median %count%', } placeholder = re.compile(r'%([^%\s]+)%') + # group-by fields copied into ECS related.* + related_users = {'user.name', 'winlog.event_data.TargetUserName', 'winlog.event_data.SubjectUserName'} + related_hosts = {'host.name', 'host.hostname', 'winlog.computer_name'} @staticmethod def lookup(doc, dotted): @@ -61,6 +65,47 @@ class SecurityOnionESAlerter(Alerter): return hashlib.sha256(key.encode('utf-8')).hexdigest() + def group(self, match): + """ Group-by values joined as ElastAlert joins them for the realert silence key. """ + return ', '.join(str(self.lookup(match, k)) for k in self.query_keys()) + + def related_bucket(self, key): + if key == 'ip' or key.endswith('.ip'): + return 'ip' + if key in self.related_users or key.endswith('.user.name'): + return 'user' + if key in self.related_hosts: + return 'hosts' + return None + + @staticmethod + def valid_ip(value): + try: + ipaddress.ip_address(value) + return True + except ValueError: + return False + + def related(self, match): + """ ECS related.* from group-by values; invalid IPs are skipped, as they fail the ip mapping. """ + related = {} + for key in self.query_keys(): + bucket = self.related_bucket(key) + if not bucket: + continue + value = self.lookup(match, key) + for v in value if isinstance(value, list) else [value]: + if v is None or (bucket == 'ip' and not self.valid_ip(str(v))): + continue + related.setdefault(bucket, {})[str(v)] = None + return {bucket: list(values) for bucket, values in related.items()} + + def event_data(self, match): + """ The match without the compound query_key field, which ES would map by its last part (e.g. .ip). """ + if not self.rule.get('compound_query_key'): + return match + return {k: v for k, v in match.items() if k != self.rule['query_key']} + @staticmethod def format_value(value): if isinstance(value, list): @@ -120,12 +165,6 @@ class SecurityOnionESAlerter(Alerter): def alert(self, matches): for match in matches: timestamp = strftime("%Y-%m-%d"'T'"%H:%M:%S"'.000Z', gmtime()) - headers = {"Content-Type": "application/json"} - - creds = None - if 'es_username' in self.rule and 'es_password' in self.rule: - creds = (self.rule['es_username'], self.rule['es_password']) - # Start building the rule dict rule_info = { "name": self.rule['detection_title'], @@ -138,39 +177,74 @@ class SecurityOnionESAlerter(Alerter): if field in self.rule: rule_info[rule_key] = self.rule[field] - summary = self.summary(match) - if summary: - rule_info["summary"] = summary + event_info = { + "kind": "alert", + "severity": self.rule['event.severity'], + "module": self.rule['event.module'], + "dataset": self.rule['event.dataset'], + "severity_label": self.rule['sigma_level'] + } + + reason = self.summary(match) + if reason: + event_info["reason"] = reason # Construct the payload with the conditional rule_info payload = { - "tags": "alert", + "tags": ["alert"], "rule": rule_info, - "event": { - "severity": self.rule['event.severity'], - "module": self.rule['event.module'], - "dataset": self.rule['event.dataset'], - "severity_label": self.rule['sigma_level'] - }, + "event": event_info, "sigma_level": self.rule['sigma_level'], - "event_data": match, + "event_data": self.event_data(match), "@timestamp": timestamp } + + keys = self.query_keys() + if keys: + payload["labels"] = { + "correlation_group_by": ', '.join(keys), + "correlation_group": self.group(match), + } + related = self.related(match) + if related: + payload["related"] = related alert_id = self.alert_id(match) # _create returns 409 on a repeat id; EAException makes ElastAlert retry url = (f"https://{self.rule['es_host']}:{self.rule['es_port']}" f"/logs-detections.alerts-so/_create/{alert_id}") - try: - response = requests.put(url, data=json.dumps(payload, cls=DateTimeEncoder), headers=headers, verify=False, auth=creds) - except requests.RequestException as e: - raise EAException(f"Unable to write alert: {e}") + response = self.put_alert(url, payload) if response.status_code == 400: - # mapping rejections fail the same way on retry, so drop them - elastalert_logger.error("Dropping alert %s for rule %s, rejected by Elasticsearch: %s", - alert_id, self.rule['detection_public_id'], response.text[:500]) - continue + # mapping rejections come from event_data; retry with it as unindexed text + rejection = response.text[:500] + payload = self.without_event_data(payload, rejection) + response = self.put_alert(url, payload) + if response.status_code == 400: + elastalert_logger.error("Dropping alert %s for rule %s, rejected by Elasticsearch even without its event data: %s; first rejection: %s", + alert_id, self.rule['detection_public_id'], response.text[:500], rejection) + continue + elastalert_logger.warning("Stored alert %s for rule %s with its event data as text, rejected by Elasticsearch: %s", + alert_id, self.rule['detection_public_id'], rejection) if response.status_code != 409 and not response.ok: raise EAException(f"Unable to write alert: {response.status_code} {response.text[:500]}") + def put_alert(self, url, payload): + creds = None + if 'es_username' in self.rule and 'es_password' in self.rule: + creds = (self.rule['es_username'], self.rule['es_password']) + try: + return requests.put(url, data=json.dumps(payload, cls=DateTimeEncoder), + headers={"Content-Type": "application/json"}, verify=False, auth=creds) + except requests.RequestException as e: + raise EAException(f"Unable to write alert: {e}") + + @staticmethod + def without_event_data(payload, rejection): + """ event_data moved to event.original; the tag keeps Fleet's final pipeline from removing it. """ + fallback = {k: v for k, v in payload.items() if k != 'event_data'} + fallback['event'] = dict(payload['event'], original=json.dumps(payload['event_data'], cls=DateTimeEncoder)) + fallback['error'] = {'message': f"event_data rejected by Elasticsearch: {rejection}"} + fallback['tags'] = payload['tags'] + ['preserve_original_event'] + return fallback + def get_info(self): return {'type': 'SecurityOnionESAlerter'} diff --git a/salt/elastalert/files/modules/so/securityonion-es_test.py b/salt/elastalert/files/modules/so/securityonion-es_test.py new file mode 100644 index 000000000..4915fec90 --- /dev/null +++ b/salt/elastalert/files/modules/so/securityonion-es_test.py @@ -0,0 +1,217 @@ +# Copyright Security Onion Solutions LLC and/or licensed to Security Onion Solutions LLC under one +# or more contributor license agreements. Licensed under the Elastic License 2.0 as shown at +# https://securityonion.net/license; you may not use this file except in compliance with the +# Elastic License 2.0. + +import copy +import importlib.util +import json +import logging +import os +import sys +import types +import unittest +from unittest.mock import MagicMock, patch + +# Real ElastAlert when installed (so-elastalert); otherwise stand-ins for what the alerter imports. +try: + import elastalert.alerts # noqa: F401 + HAVE_ELASTALERT = True +except ImportError: + HAVE_ELASTALERT = False + + class Alerter: + def __init__(self, rule): + self.rule = rule + + class DateTimeEncoder(json.JSONEncoder): + def default(self, obj): + return obj.isoformat() if hasattr(obj, 'isoformat') else json.JSONEncoder.default(self, obj) + + class EAException(Exception): + pass + + alerts = types.ModuleType('elastalert.alerts') + alerts.Alerter = Alerter + alerts.DateTimeEncoder = DateTimeEncoder + util = types.ModuleType('elastalert.util') + util.EAException = EAException + util.elastalert_logger = logging.getLogger('elastalert') + package = types.ModuleType('elastalert') + package.alerts = alerts + package.util = util + sys.modules.update({'elastalert': package, 'elastalert.alerts': alerts, 'elastalert.util': util}) + +spec = importlib.util.spec_from_file_location('securityonion_es', os.path.join(os.path.dirname(__file__), 'securityonion-es.py')) +es = importlib.util.module_from_spec(spec) +spec.loader.exec_module(es) + +BASE_RULE = { + 'name': 'Many Failed Network Logons To One Host From One Source -- 35a42db6-6629-45af-b8aa-e1fa33c28ef5', + 'detection_title': 'Many Failed Network Logons To One Host From One Source', + 'detection_public_id': '35a42db6-6629-45af-b8aa-e1fa33c28ef5', + 'sigma_level': 'medium', + 'sigma_correlation': 'event_count', + 'event.severity': 3, + 'event.module': 'sigma', + 'event.dataset': 'sigma.alert', + 'es_host': 'manager', + 'es_port': 9200, + 'summary_template': '%count% failed network logons to %host.name% from %source.ip% in %duration%', +} + + +def correlation_match(): + return { + 'event_count': 3561, + 'window_start': '2026-09-30T18:05:10+00:00', + '@timestamp': '2026-09-30T18:07:53+00:00', + 'host': {'name': 'sa-delta-02-jb'}, + 'source': {'ip': '192.168.198.149'}, + '_id': '6d1c', + 'num_hits': 1, + 'num_matches': 1, + } + + +class TestSecurityOnionESAlerter(unittest.TestCase): + + def send(self, rule, match): + """ Run alert() and return the payload it wrote and the URL it wrote to. """ + alerter = es.SecurityOnionESAlerter(rule) + response = MagicMock(status_code=201, ok=True) + with patch.object(es.requests, 'put', return_value=response) as put: + alerter.alert([match]) + self.assertEqual(put.call_count, 1) + return json.loads(put.call_args.kwargs['data']), put.call_args.args[0] + + def test_compound_query_key_left_out_of_event_data(self): + rule = dict(BASE_RULE, compound_query_key=['host.name', 'source.ip'], query_key='host.name,source.ip') + match = correlation_match() + match['host.name,source.ip'] = 'sa-delta-02-jb, 192.168.198.149' + original = copy.deepcopy(match) + + payload, url = self.send(rule, match) + + self.assertNotIn('host.name,source.ip', payload['event_data']) + self.assertEqual(payload['event_data']['host'], {'name': 'sa-delta-02-jb'}) + self.assertEqual(payload['event_data']['source'], {'ip': '192.168.198.149'}) + self.assertEqual(payload['labels'], {'correlation_group_by': 'host.name, source.ip', 'correlation_group': 'sa-delta-02-jb, 192.168.198.149'}) + self.assertEqual(payload['related'], {'hosts': ['sa-delta-02-jb'], 'ip': ['192.168.198.149']}) + self.assertEqual(payload['event']['kind'], 'alert') + self.assertEqual(payload['event']['reason'], '3,561 failed network logons to sa-delta-02-jb from 192.168.198.149 in 2 minutes') + self.assertNotIn('summary', payload['rule']) + # ElastAlert reuses the match + self.assertEqual(match, original) + # id from the group fields, not the compound key + without_key = {k: v for k, v in match.items() if k != 'host.name,source.ip'} + self.assertTrue(url.endswith('/' + es.SecurityOnionESAlerter(rule).alert_id(without_key))) + + def test_single_query_key(self): + rule = dict(BASE_RULE, query_key='user.name', summary_template='%count% failed SOC logins for %user.name%') + match = {'event_count': 3, 'window_start': '2026-09-30T16:49:52+00:00', '@timestamp': '2026-09-30T16:50:00+00:00', + 'user': {'name': 'josh@local.invalid'}} + + payload, _ = self.send(rule, match) + + self.assertEqual(payload['labels'], {'correlation_group_by': 'user.name', 'correlation_group': 'josh@local.invalid'}) + self.assertEqual(payload['related'], {'user': ['josh@local.invalid']}) + self.assertEqual(payload['event']['reason'], '3 failed SOC logins for josh@local.invalid') + self.assertEqual(payload['event_data'], match) + + def test_related_buckets(self): + rule = dict(BASE_RULE, compound_query_key=['winlog.event_data.TargetUserName', 'source.ip', 'dns.highest_registered_domain'], + query_key='winlog.event_data.TargetUserName,source.ip,dns.highest_registered_domain') + match = correlation_match() + match.update({'winlog': {'event_data': {'TargetUserName': ['admmig', 'svc', 'admmig']}}, + 'source': {'ip': 'not-an-ip'}, 'dns': {'highest_registered_domain': 'example.com'}}) + + payload, _ = self.send(rule, match) + + # deduped; invalid IP skipped; domain stays in the group only + self.assertEqual(payload['related'], {'user': ['admmig', 'svc']}) + self.assertEqual(payload['labels']['correlation_group'], "['admmig', 'svc', 'admmig'], not-an-ip, example.com") + + def test_group_without_related_fields(self): + rule = dict(BASE_RULE, query_key='dns.highest_registered_domain') + match = dict(correlation_match(), dns={'highest_registered_domain': 'example.com'}) + + payload, _ = self.send(rule, match) + + self.assertEqual(payload['labels'], {'correlation_group_by': 'dns.highest_registered_domain', 'correlation_group': 'example.com'}) + self.assertNotIn('related', payload) + + def test_plain_rule_has_no_correlation_fields(self): + match = {'@timestamp': '2026-09-30T16:50:00+00:00', '_id': 'abc', 'process': {'name': 'whoami.exe'}} + + payload, url = self.send(dict(BASE_RULE), match) + + self.assertNotIn('labels', payload) + self.assertNotIn('related', payload) + self.assertNotIn('reason', payload['event']) + self.assertEqual(payload['event']['kind'], 'alert') + self.assertEqual(payload['event_data'], match) + self.assertTrue(url.endswith('/' + es.SecurityOnionESAlerter(dict(BASE_RULE)).alert_id(match))) + + def send_responses(self, rule, match, responses): + """ Run alert() against a sequence of write responses; return the payloads written. """ + alerter = es.SecurityOnionESAlerter(rule) + with patch.object(es.requests, 'put', side_effect=responses) as put: + alerter.alert([match]) + return [json.loads(c.kwargs['data']) for c in put.call_args_list] + + def test_rejected_event_data_is_kept_as_text(self): + rule = dict(BASE_RULE, query_key='user.name', summary_template='%count% failed SOC logins for %user.name%') + match = {'event_count': 3, 'window_start': '2026-09-30T16:49:52+00:00', '@timestamp': '2026-09-30T16:50:00+00:00', + 'user': {'name': 'josh@local.invalid'}} + rejected = MagicMock(status_code=400, ok=False, text='{"error":{"type":"document_parsing_exception"}}') + + first, second = self.send_responses(rule, match, [rejected, MagicMock(status_code=201, ok=True)]) + + self.assertIn('event_data', first) + self.assertNotIn('event_data', second) + self.assertEqual(json.loads(second['event']['original']), match) + self.assertEqual(second['tags'], ['alert', 'preserve_original_event']) + self.assertEqual(first['tags'], ['alert']) + self.assertTrue(second['error']['message'].startswith('event_data rejected by Elasticsearch: {"error"')) + self.assertEqual(second['event']['reason'], first['event']['reason']) + self.assertEqual(second['labels'], first['labels']) + self.assertEqual(second['related'], first['related']) + self.assertEqual(second['rule'], first['rule']) + + def test_rejected_twice_is_dropped_without_retry(self): + rejected = MagicMock(status_code=400, ok=False, text='{"error":{"type":"document_parsing_exception"}}') + match = {'@timestamp': '2026-09-30T16:50:00+00:00', '_id': 'abc'} + + # no EAException, so no retry + payloads = self.send_responses(dict(BASE_RULE), match, [rejected, rejected]) + + self.assertEqual(len(payloads), 2) + + @unittest.skipUnless(HAVE_ELASTALERT, 'needs ElastAlert, as in the so-elastalert container') + def test_group_matches_elastalert_silence_key(self): + from elastalert.elastalert import ElastAlerter + from elastalert.util import ts_to_dt + + cases = [ + (['host.name', 'source.ip'], {'host': {'name': 'sa-delta-02-jb'}, 'source': {'ip': '192.168.198.149'}}), + (['winlog.event_data.TargetUserName', 'source.ip'], {'winlog': {'event_data': {'TargetUserName': ['admmig', 'svc']}}, 'source': {'ip': '10.23.23.9'}}), + (['user.name'], {'user': {'name': "o'brien \"q\" \\ *:?@local.invalid"}}), + ] + for keys, fields in cases: + with self.subTest(keys=keys): + rule = dict(BASE_RULE, timestamp_field='@timestamp', ts_to_dt=ts_to_dt) + if len(keys) > 1: + rule.update(compound_query_key=keys, query_key=','.join(keys)) + else: + rule['query_key'] = keys[0] + hit = {'_id': 'x', '_source': dict(correlation_match(), **copy.deepcopy(fields))} + match = ElastAlerter.process_hits(rule, [hit])[0] + + silence_suffix = ElastAlerter.get_named_key_value(None, rule, match, 'query_key') + + self.assertEqual(es.SecurityOnionESAlerter(rule).group(match), silence_suffix) + + +if __name__ == '__main__': + unittest.main() diff --git a/salt/elasticsearch/files/ingest/logs-system.auth@custom b/salt/elasticsearch/files/ingest/logs-system.auth@custom new file mode 100644 index 000000000..19e3b9a3d --- /dev/null +++ b/salt/elasticsearch/files/ingest/logs-system.auth@custom @@ -0,0 +1,34 @@ +{ + "version": 1, + "_meta": { + "managed_by": "securityonion", + "managed": true + }, + "description": "Custom pipeline for the System integration's auth data stream.", + "processors": [ + { + "trim": { + "description": "Grok leaves a leading space on 'invalid user' names (elastic/integrations#12174) and, before 2.23.2, sudo padding", + "field": "user.name", + "ignore_missing": true, + "ignore_failure": true + } + }, + { + "trim": { + "description": "Appended from the untrimmed user.name", + "field": "related.user", + "ignore_missing": true, + "ignore_failure": true + } + }, + { + "script": { + "description": "Dedupe after trimming", + "if": "ctx.related?.user instanceof List", + "source": "ctx.related.user = new ArrayList(new LinkedHashSet(ctx.related.user));", + "ignore_failure": true + } + } + ] +} diff --git a/salt/elasticsearch/templates/component/so/detections-alerts-mappings.json b/salt/elasticsearch/templates/component/so/detections-alerts-mappings.json index e5b536ed9..5a2abfb7f 100644 --- a/salt/elasticsearch/templates/component/so/detections-alerts-mappings.json +++ b/salt/elasticsearch/templates/component/so/detections-alerts-mappings.json @@ -35,15 +35,6 @@ "correlation": { "ignore_above": 1024, "type": "keyword" - }, - "summary": { - "type": "match_only_text", - "fields": { - "keyword": { - "ignore_above": 1024, - "type": "keyword" - } - } } } }, @@ -63,6 +54,24 @@ "dataset": { "ignore_above": 1024, "type": "keyword" + }, + "kind": { + "ignore_above": 1024, + "type": "keyword" + }, + "reason": { + "type": "match_only_text", + "fields": { + "keyword": { + "ignore_above": 1024, + "type": "keyword" + } + } + }, + "original": { + "type": "keyword", + "index": false, + "doc_values": false } } }, @@ -96,6 +105,40 @@ "type": "double" } } + }, + "labels": { + "properties": { + "correlation_group_by": { + "ignore_above": 1024, + "type": "keyword" + }, + "correlation_group": { + "ignore_above": 1024, + "type": "keyword" + } + } + }, + "related": { + "properties": { + "ip": { + "type": "ip" + }, + "user": { + "ignore_above": 1024, + "type": "keyword" + }, + "hosts": { + "ignore_above": 1024, + "type": "keyword" + } + } + }, + "error": { + "properties": { + "message": { + "type": "match_only_text" + } + } } } } diff --git a/salt/soc/defaults.yaml b/salt/soc/defaults.yaml index 1296fe48e..04ef707af 100644 --- a/salt/soc/defaults.yaml +++ b/salt/soc/defaults.yaml @@ -1445,7 +1445,7 @@ soc: default: - repo: https://github.com/Security-Onion-Solutions/securityonion-resources license: Elastic-2.0 - folder: sigma/stable + folder: sigma community: true rulesetName: securityonion-resources - repo: file:///nsm/rules/custom-local-repos/local-sigma @@ -1455,7 +1455,7 @@ soc: airgap: - repo: file:///nsm/rules/detect-sigma/repos/securityonion-resources license: Elastic-2.0 - folder: sigma/stable + folder: sigma community: true rulesetName: securityonion-resources - repo: file:///nsm/rules/custom-local-repos/local-sigma @@ -1465,7 +1465,7 @@ soc: sigmaRulePackages: - core - emerging_threats_addon - useEsql: false + useEsql: true esqlCaseInsensitive: true esqlQueryDelaySeconds: 30 esqlCorrelationAllowanceSeconds: 600 diff --git a/salt/soc/soc_soc.yaml b/salt/soc/soc_soc.yaml index c6eab6c1d..2aa54627b 100644 --- a/salt/soc/soc_soc.yaml +++ b/salt/soc/soc_soc.yaml @@ -400,7 +400,7 @@ soc: advanced: False helpLink: sigma useEsql: - description: "(Pre-release) Use Elasticsearch Piped Query Language (ES|QL) instead of EQL (Elastic Query Language) for Elasticsearch queries. The Sigma converter will output ES|QL instead of EQL, allowing support for correlations." + description: "(Pre-release) Use Elasticsearch Piped Query Language (ES|QL) instead of EQL (Elastic Query Language) for Elasticsearch queries. The Sigma converter will output ES|QL instead of EQL, allowing support for correlations. Switching back to EQL is not supported for correlations: they stay enabled but stop running until ES|QL is turned on again." global: True advanced: True forcedType: bool diff --git a/setup/so-verify b/setup/so-verify index 660424c72..0cd910021 100755 --- a/setup/so-verify +++ b/setup/so-verify @@ -67,7 +67,7 @@ log_has_errors() { grep -vE "Reading first line of patchfile" | \ grep -vE "Command failed with exit code" | \ grep -vE "Running scope as unit" | \ - grep -vE "securityonion-resources/sigma/stable" | \ + grep -vE "securityonion-resources/sigma/" | \ grep -vE "remove_failed_vm.sls" | \ grep -vE "failed to copy: httpReadSeeker" | \ grep -vE "Error response from daemon: failed to resolve reference" | \