diff --git a/.github/workflows/pythontest.yml b/.github/workflows/pythontest.yml index 5d474f5df..7e5c2742e 100644 --- a/.github/workflows/pythontest.yml +++ b/.github/workflows/pythontest.yml @@ -6,6 +6,7 @@ on: - "salt/sensoroni/files/analyzers/**" - "salt/manager/tools/sbin/**" - "salt/_beacons/**" + - "salt/elastalert/files/modules/**" - "salt/telegraf/tools/sbin_jinja/**" - "salt/telegraf/defaults.yaml" - "salt/telegraf/soc_telegraf.yaml" @@ -18,7 +19,7 @@ jobs: fail-fast: false matrix: python-version: ["3.14"] - python-code-path: ["salt/sensoroni/files/analyzers", "salt/manager/tools/sbin", "salt/_beacons"] + python-code-path: ["salt/sensoroni/files/analyzers", "salt/manager/tools/sbin", "salt/_beacons", "salt/elastalert/files/modules/so"] steps: - uses: actions/checkout@v3 diff --git a/salt/elastalert/defaults.yaml b/salt/elastalert/defaults.yaml index 393932992..1986fdce3 100644 --- a/salt/elastalert/defaults.yaml +++ b/salt/elastalert/defaults.yaml @@ -10,7 +10,7 @@ elastalert: buffer_time: minutes: 10 old_query_limit: - minutes: 5 + minutes: 1440 es_port: 9200 es_conn_timeout: 55 max_query_size: 5000 diff --git a/salt/elastalert/files/modules/so/conftest.py b/salt/elastalert/files/modules/so/conftest.py new file mode 100644 index 000000000..2220e0a7e --- /dev/null +++ b/salt/elastalert/files/modules/so/conftest.py @@ -0,0 +1,80 @@ +# 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. + +from datetime import datetime +import json +import logging +import sys +import types + +# stand-ins when ElastAlert isn't installed (CI) +try: + import elastalert.alerts # noqa: F401 +except ImportError: + + 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 + + def lookup_es_key(doc, term): + for part in term.split('.'): + if not isinstance(doc, dict) or part not in doc: + return None + doc = doc[part] + return doc + + def ts_to_dt(value): + return value if isinstance(value, datetime) else datetime.fromisoformat(value) + + def elasticsearch_client(conf): + return None # tests set the alerter's client + + 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') + util.lookup_es_key = lookup_es_key + util.ts_to_dt = ts_to_dt + util.elasticsearch_client = elasticsearch_client + package = types.ModuleType('elastalert') + package.alerts = alerts + package.util = util + sys.modules.update({'elastalert': package, 'elastalert.alerts': alerts, 'elastalert.util': util}) + +# stand-ins when elasticsearch-py isn't installed (CI) +try: + import elasticsearch.exceptions # noqa: F401 +except ImportError: + + class ElasticsearchException(Exception): + pass + + class TransportError(ElasticsearchException): + pass + + class ConnectionError(TransportError): + pass + + class ConflictError(TransportError): + pass + + class RequestError(TransportError): + pass + + exceptions = types.ModuleType('elasticsearch.exceptions') + for cls in (ElasticsearchException, TransportError, ConnectionError, ConflictError, RequestError): + setattr(exceptions, cls.__name__, cls) + es_package = types.ModuleType('elasticsearch') + es_package.exceptions = exceptions + sys.modules.update({'elasticsearch': es_package, 'elasticsearch.exceptions': exceptions}) diff --git a/salt/elastalert/files/modules/so/securityonion-es.py b/salt/elastalert/files/modules/so/securityonion-es.py index d9bb8009e..8104661e5 100644 --- a/salt/elastalert/files/modules/so/securityonion-es.py +++ b/salt/elastalert/files/modules/so/securityonion-es.py @@ -1,63 +1,269 @@ # -*- coding: utf-8 -*- # 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 +# 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. - -from time import gmtime, strftime -import requests,json -from elastalert.alerts import Alerter +from datetime import datetime, timezone +import hashlib +import ipaddress +import json +import re +import uuid import urllib3 +from elasticsearch.exceptions import ConflictError, ElasticsearchException, RequestError +from elastalert.alerts import Alerter, DateTimeEncoder +from elastalert.util import EAException, elastalert_logger, elasticsearch_client, lookup_es_key, ts_to_dt + +# grid runs verify_certs: false; also quiets ElastAlert's own queries urllib3.disable_warnings(urllib3.exceptions.InsecureRequestWarning) +ALERT_INDEX = 'logs-detections.alerts-so' +# ES error text kept in logs and alerts +ERROR_TEXT_LIMIT = 500 +# a match missing backend columns (window_start, count, @timestamp); the alert is still written +MATCH_ERRORS = (KeyError, TypeError, ValueError) + + class SecurityOnionESAlerter(Alerter): """ Use matched data to create alerts in Elasticsearch. """ - required_options = set(['detection_title', 'sigma_level']) - optional_fields = ['sigma_category', 'sigma_product', 'sigma_service'] + required_options = {'detection_title', 'sigma_level'} + optional_fields = ['sigma_category', 'sigma_product', 'sigma_service', 'sigma_correlation'] + + # count column and default summary per type; stored alert data, so not localized + CORRELATION_COUNTS = { + 'event_count': ('event_count', '%count% events'), + 'value_count': ('value_count', '%count% distinct values'), + 'temporal': ('event_type_count', '%count% correlated rules matched'), + 'value_sum': ('value_sum', 'total %count%'), + 'value_avg': ('value_avg', 'average %count%'), + 'value_percentile': ('value_percentile', 'percentile %count%'), + 'value_median': ('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'} + + def __init__(self, rule): + super().__init__(rule) + # uses the grid's TLS, auth and timeout settings + self.es = elasticsearch_client(rule) + + @property + def is_correlation(self): + return bool(self.rule.get('sigma_correlation')) + + def query_keys(self): + """compound_query_key holds the list; query_key is flattened to a string.""" + if self.rule.get('compound_query_key'): + return self.rule['compound_query_key'] + if self.rule.get('query_key'): + return [self.rule['query_key']] + return [] + + def alert_id(self, match): + """Stable id: window end + group values for correlations, source _id otherwise; random without one.""" + if self.is_correlation: + # ungrouped rows have a hashed _id that changes with the count + values = ''.join(f"|{lookup_es_key(match, k)}" for k in self.query_keys()) + key = f"{self.rule['detection_public_id']}|{ts_to_dt(match['@timestamp']).isoformat()}{values}" + elif match.get('_id'): + key = f"{self.rule['detection_public_id']}|{match['_id']}" + else: + return uuid.uuid4().hex + + return hashlib.sha256(key.encode('utf-8')).hexdigest() + + def group(self, match): + """Group-by values, joined like ElastAlert's realert key.""" + return ', '.join(str(lookup_es_key(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; skips invalid IPs.""" + related = {} + for key in self.query_keys(): + bucket = self.related_bucket(key) + if not bucket: + continue + # original spellings of a lowercased group + value = lookup_es_key(match, f"{key}_spellings") + if value is None: + value = lookup_es_key(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 + # dict: ordered and deduped + related.setdefault(bucket, {})[str(v)] = None + return {bucket: list(values) for bucket, values in related.items()} + + def event_data(self, match): + """The match minus the compound query_key field, which ES would map by its last part.""" + 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): + shown = ', '.join(str(v) for v in value[:3]) + return shown if len(value) <= 3 else f"{shown} and {len(value) - 3} more" + return str(value) + + @staticmethod + def format_count(value): + if isinstance(value, float) and not value.is_integer(): + return f"{value:,.2f}" + if isinstance(value, (int, float)): + return f"{int(value):,}" + return str(value) + + @staticmethod + def format_duration(seconds): + for unit, size in (('hour', 3600), ('minute', 60)): + if seconds >= 2 * size: + return f"{seconds // size} {unit}s" + return f"{seconds} second{'' if seconds == 1 else 's'}" + + def summary(self, match): + """One-line correlation summary.""" + column, label = self.CORRELATION_COUNTS.get(self.rule['sigma_correlation'], (None, '%count%')) + start = ts_to_dt(match['window_start']) + end = ts_to_dt(match['@timestamp']) + values = { + 'count': self.format_count(match.get(column)), + 'start': start.strftime('%Y-%m-%d %H:%M:%S UTC'), + 'end': end.strftime('%Y-%m-%d %H:%M:%S UTC'), + 'duration': self.format_duration(int((end - start).total_seconds())), + } + + template = self.rule.get('summary_template') + if not template: + groups = ', '.join(f"{k} %{k}%" for k in self.query_keys()) + template = f"{label} for {groups} in %duration%" if groups else f"{label} in %duration%" + + def fill(m): + if m[1] in values: + return values[m[1]] + value = lookup_es_key(match, m[1]) + # unknown placeholders stay visible so typos show + return m[0] if value is None else self.format_value(value) + + return self.PLACEHOLDER.sub(fill, template) def alert(self, matches): for match in matches: - timestamp = strftime("%Y-%m-%d"'T'"%H:%M:%S"'.000Z', gmtime()) - headers = {"Content-Type": "application/json"} + try: + alert_id = self.alert_id(match) + except MATCH_ERRORS as e: + elastalert_logger.warning("Writing alert for rule %s without a stable id, so a retry may duplicate it: %r", + self.rule['detection_public_id'], e) + alert_id = uuid.uuid4().hex + try: + self.write(alert_id, self.payload(match)) + except ElasticsearchException as e: + # EAException makes ElastAlert retry + raise EAException(f"Unable to write the alert to Elasticsearch: {str(e)[:ERROR_TEXT_LIMIT]}") from e - creds = None - if 'es_username' in self.rule and 'es_password' in self.rule: - creds = (self.rule['es_username'], self.rule['es_password']) + def payload(self, match): + rule_info = { + "name": self.rule['detection_title'], + "uuid": self.rule['detection_public_id'] + } - # Start building the rule dict - rule_info = { - "name": self.rule['detection_title'], - "uuid": self.rule['detection_public_id'] - } + # Add optional fields if they are present in the rule + for field in self.optional_fields: + rule_key = field.split('_')[-1] # Assumes field format "sigma_" + if field in self.rule: + rule_info[rule_key] = self.rule[field] - # Add optional fields if they are present in the rule - for field in self.optional_fields: - rule_key = field.split('_')[-1] # Assumes field format "sigma_" - if field in self.rule: - rule_info[rule_key] = self.rule[field] + 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'] + } - # Construct the payload with the conditional rule_info - payload = { - "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'] - }, - "sigma_level": self.rule['sigma_level'], - "event_data": match, - "@timestamp": timestamp - } - url = f"https://{self.rule['es_host']}:{self.rule['es_port']}/logs-detections.alerts-so/_doc/" - requests.post(url, data=json.dumps(payload), headers=headers, verify=False, auth=creds) + payload = { + "tags": ["alert"], + "rule": rule_info, + "event": event_info, + "sigma_level": self.rule['sigma_level'], + "event_data": self.event_data(match), + "@timestamp": datetime.now(timezone.utc).strftime('%Y-%m-%dT%H:%M:%S.000Z') + } + + if self.is_correlation: + keys = self.query_keys() + try: + # built before any is added, so a failure adds none + reason = self.summary(match) + labels = {"correlation_group_by": ', '.join(keys), "correlation_group": self.group(match)} if keys else None + related = self.related(match) + except MATCH_ERRORS as e: + elastalert_logger.warning("Writing alert for rule %s without its correlation summary: %r", + self.rule['detection_public_id'], e) + else: + payload["event"]["reason"] = reason + if labels: + payload["labels"] = labels + if related: + payload["related"] = related + + return payload + + def write(self, alert_id, payload): + try: + self.create(alert_id, payload) + except RequestError as e: + # mapping rejections come from event_data; retry it as text + rejection = str(e)[:ERROR_TEXT_LIMIT] + try: + self.create(alert_id, self.without_event_data(payload, rejection)) + except RequestError as again: + 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'], str(again)[:ERROR_TEXT_LIMIT], rejection) + return + 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) + + def create(self, alert_id, payload): + try: + self.es.create(index=ALERT_INDEX, id=alert_id, body=payload) + except ConflictError: + pass # a repeat id is already stored + + @staticmethod + def without_event_data(payload, rejection): + """Moves event_data to event.original; the tag keeps Fleet's final pipeline from dropping 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..8f4db168d --- /dev/null +++ b/salt/elastalert/files/modules/so/securityonion-es_test.py @@ -0,0 +1,250 @@ +# 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 +from datetime import datetime, timezone +import importlib.util +import json +import os +import unittest +from unittest.mock import MagicMock + +from elasticsearch.exceptions import ConflictError, ConnectionError, RequestError + +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, + 'es_conn_timeout': 55, + 'summary_template': '%count% failed network logons to %host.name% from %source.ip% in %duration%', +} + +PLAIN_RULE = {k: v for k, v in BASE_RULE.items() if k not in ('sigma_correlation', 'summary_template')} + + +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': 'host-01'}, + 'source': {'ip': '192.0.2.10'}, + '_id': '6d1c', + 'num_hits': 1, + 'num_matches': 1, + } + + +class TestSecurityOnionESAlerter(unittest.TestCase): + + def creates(self, rule, match, effects=None): + """Run alert(); return (body, id) of each create.""" + alerter = es.SecurityOnionESAlerter(rule) + alerter.es = MagicMock() + alerter.es.create.side_effect = effects + alerter.alert([match]) + calls = alerter.es.create.call_args_list + self.assertTrue(all(c.kwargs['index'] == 'logs-detections.alerts-so' for c in calls)) + # as the client serializes it + return [(json.loads(json.dumps(c.kwargs['body'], cls=es.DateTimeEncoder)), c.kwargs['id']) for c in calls] + + def send(self, rule, match): + """Run alert(); return the payload it wrote and its id.""" + (payload, alert_id), = self.creates(rule, match) + return payload, alert_id + + 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'] = 'host-01, 192.0.2.10' + original = copy.deepcopy(match) + + payload, _ = self.send(rule, match) + + self.assertNotIn('host.name,source.ip', payload['event_data']) + self.assertEqual(payload['event_data']['host'], {'name': 'host-01'}) + self.assertEqual(payload['event_data']['source'], {'ip': '192.0.2.10'}) + self.assertEqual(payload['labels'], {'correlation_group_by': 'host.name, source.ip', 'correlation_group': 'host-01, 192.0.2.10'}) + self.assertEqual(payload['related'], {'hosts': ['host-01'], 'ip': ['192.0.2.10']}) + self.assertEqual(payload['event']['kind'], 'alert') + self.assertEqual(payload['event']['reason'], '3,561 failed network logons to host-01 from 192.0.2.10 in 2 minutes') + # ElastAlert reuses the match + self.assertEqual(match, original) + + 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': 'user@example.invalid'}} + + payload, _ = self.send(rule, match) + + self.assertEqual(payload['labels'], {'correlation_group_by': 'user.name', 'correlation_group': 'user@example.invalid'}) + self.assertEqual(payload['related'], {'user': ['user@example.invalid']}) + self.assertEqual(payload['event']['reason'], '3 failed SOC logins for user@example.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': ['admin1', 'svc', 'admin1']}}, + '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': ['admin1', 'svc']}) + self.assertEqual(payload['labels']['correlation_group'], "['admin1', 'svc', 'admin1'], not-an-ip, example.com") + + payload, _ = self.send(dict(BASE_RULE, query_key='dns.highest_registered_domain'), match) + + self.assertNotIn('related', payload) + + def test_related_uses_original_spellings(self): + rule = dict(BASE_RULE, query_key='user.name', summary_template=None) + match = {'event_count': 3, 'window_start': '2026-09-30T16:49:52+00:00', '@timestamp': '2026-09-30T16:50:00+00:00', + 'user': {'name': 'admin', 'name_spellings': ['Admin', 'admin', 'ADMIN']}} + + payload, _ = self.send(rule, match) + + # the group shows the lowercased value; related.user finds every spelling + self.assertEqual(payload['labels']['correlation_group'], 'admin') + self.assertEqual(payload['related'], {'user': ['Admin', 'admin', 'ADMIN']}) + self.assertEqual(payload['event']['reason'], '3 events for user.name admin in 8 seconds') + + def test_plain_rule_has_no_correlation_fields(self): + # a query_key alone (e.g. from an override) isn't a correlation + rule = dict(PLAIN_RULE, query_key='user.name') + # ElastAlert parses @timestamp for EQL hits + match = {'@timestamp': datetime(2026, 9, 30, 16, 50, tzinfo=timezone.utc), '_id': 'abc', + 'process': {'name': 'whoami.exe'}, 'user': {'name': 'user'}} + + payload, alert_id = self.send(rule, match) + + alerter = es.SecurityOnionESAlerter(rule) + self.assertNotIn('labels', payload) + self.assertNotIn('related', payload) + self.assertNotIn('reason', payload['event']) + self.assertEqual(payload['event']['kind'], 'alert') + self.assertEqual(payload['event_data'], dict(match, **{'@timestamp': '2026-09-30T16:50:00+00:00'})) + self.assertEqual(alert_id, alerter.alert_id(match)) + self.assertNotEqual(alerter.alert_id(match), alerter.alert_id(dict(match, _id='abd'))) + # without an _id, never deduplicated + self.assertNotEqual(alerter.alert_id({'_id': None}), alerter.alert_id({'_id': None})) + + def test_ungrouped_correlation_id_ignores_row_hash(self): + rule = dict(BASE_RULE, summary_template=None) + first = {k: v for k, v in correlation_match().items() if k not in ('host', 'source')} + # ES|QL hashes the row into _id, so a later count changes it + later = dict(first, event_count=3600, _id='9f2a') + + payload, alert_id = self.send(rule, first) + + alerter = es.SecurityOnionESAlerter(rule) + self.assertEqual(alerter.alert_id(first), alerter.alert_id(later)) + self.assertNotEqual(alerter.alert_id(first), alerter.alert_id(dict(first, **{'@timestamp': '2026-09-30T18:09:00+00:00'}))) + self.assertEqual(alert_id, alerter.alert_id(first)) + self.assertEqual(payload['event']['reason'], '3,561 events in 2 minutes') + self.assertNotIn('labels', payload) + + def test_temporal_count_column(self): + rule = dict(BASE_RULE, sigma_correlation='temporal', summary_template=None) + match = {'event_type_count': 2, 'window_start': '2026-09-30T16:49:52+00:00', '@timestamp': '2026-09-30T16:50:00+00:00'} + + payload, _ = self.send(rule, match) + + self.assertEqual(payload['event']['reason'], '2 correlated rules matched in 8 seconds') + + def test_grouped_correlation_id_unchanged(self): + """Ids of alerts already written must not change.""" + rule = dict(BASE_RULE, compound_query_key=['host.name', 'source.ip'], query_key='host.name,source.ip') + key = f"{BASE_RULE['detection_public_id']}|2026-09-30T18:07:53+00:00|host-01|192.0.2.10" + + self.assertEqual(es.SecurityOnionESAlerter(rule).alert_id(correlation_match()), es.hashlib.sha256(key.encode()).hexdigest()) + + 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': 'user@example.invalid'}} + rejected = RequestError(400, 'document_parsing_exception', {}) + + (first, _), (second, _) = self.creates(rule, match, [rejected, None]) + + 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: ')) + self.assertIn('document_parsing_exception', second['error']['message']) + # everything else carries over + self.assertEqual({k: v for k, v in second['event'].items() if k != 'original'}, first['event']) + changed = ('event_data', 'event', 'error', 'tags') + self.assertEqual({k: v for k, v in second.items() if k not in changed}, {k: v for k, v in first.items() if k not in changed}) + + def test_rejected_twice_is_dropped_without_retry(self): + rejected = RequestError(400, 'document_parsing_exception', {}) + match = {'@timestamp': '2026-09-30T16:50:00+00:00', '_id': 'abc'} + + # no EAException, so no retry + self.assertEqual(len(self.creates(PLAIN_RULE, match, [rejected, rejected])), 2) + + def test_write_failure_is_retried(self): + match = {'@timestamp': '2026-09-30T16:50:00+00:00', '_id': 'abc'} + + # EAException makes ElastAlert retry the alert + with self.assertRaisesRegex(es.EAException, 'Unable to write the alert to Elasticsearch'): + self.creates(PLAIN_RULE, match, [ConnectionError('N/A', 'refused', None)]) + + # a repeat id is already stored + self.assertEqual(len(self.creates(PLAIN_RULE, match, [ConflictError(409, 'version_conflict_engine_exception', {})])), 1) + + def test_correlation_fields_are_optional(self): + rule = dict(BASE_RULE, query_key='source.ip') + # no window_start: the summary cannot be built + match = {k: v for k, v in correlation_match().items() if k != 'window_start'} + + with self.assertLogs('elastalert', 'WARNING'): + payload, _ = self.send(rule, match) + + self.assertEqual(payload['event']['kind'], 'alert') + self.assertNotIn('reason', payload['event']) + self.assertNotIn('labels', payload) + self.assertNotIn('related', payload) + + def test_unstable_id_still_writes(self): + # no @timestamp: the window end is unknown + match = {k: v for k, v in correlation_match().items() if k not in ('@timestamp', 'window_start')} + + with self.assertLogs('elastalert', 'WARNING'): + payload, alert_id = self.send(BASE_RULE, match) + + self.assertEqual(len(alert_id), 32) + self.assertEqual(payload['event_data'], match) + + def test_summary_formats_values(self): + rule = dict(BASE_RULE, sigma_correlation='value_avg', query_key='source.ip', + summary_template='%count% for %source.ip% to %destination.port%') + match = {'value_avg': 2.5, 'window_start': '2026-09-30T16:49:52+00:00', '@timestamp': '2026-09-30T16:50:00+00:00', + 'source': {'ip': '192.0.2.10'}, 'destination': {'port': [22, 80, 443, 8080, 8443]}} + + payload, _ = self.send(rule, match) + + self.assertEqual(payload['event']['reason'], '2.50 for 192.0.2.10 to 22, 80, 443 and 2 more') + self.assertEqual(es.SecurityOnionESAlerter.format_count('n/a'), 'n/a') + + def test_get_info(self): + self.assertEqual(es.SecurityOnionESAlerter(PLAIN_RULE).get_info(), {'type': 'SecurityOnionESAlerter'}) diff --git a/salt/elastalert/soc_elastalert.yaml b/salt/elastalert/soc_elastalert.yaml index 123ead697..3d848e5af 100644 --- a/salt/elastalert/soc_elastalert.yaml +++ b/salt/elastalert/soc_elastalert.yaml @@ -120,7 +120,7 @@ elastalert: helpLink: elastalert old_query_limit: minutes: - description: Amount of time in minutes between queries to start at the most recently run query. + description: How long ElastAlert can be down, in minutes, and still resume each rule where it stopped. After a longer outage, rules restart from now and skip the gap. global: True helpLink: elastalert es_conn_timeout: diff --git a/salt/elasticsearch/defaults.yaml b/salt/elasticsearch/defaults.yaml index c8ceab34d..ab313a13f 100644 --- a/salt/elasticsearch/defaults.yaml +++ b/salt/elasticsearch/defaults.yaml @@ -1160,6 +1160,7 @@ elasticsearch: - so-fleet_agent_id_verification-1 - so-logs-mappings - so-logs-settings + - detections-alerts-mappings data_stream: allow_custom_routing: false hidden: false 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/detection-mappings.json b/salt/elasticsearch/templates/component/so/detection-mappings.json index 4dd5b45e7..1098e440a 100644 --- a/salt/elasticsearch/templates/component/so/detection-mappings.json +++ b/salt/elasticsearch/templates/component/so/detection-mappings.json @@ -50,6 +50,18 @@ "ignore_above": 1024, "type": "keyword" }, + "ruleType": { + "ignore_above": 1024, + "type": "keyword" + }, + "correlationType": { + "ignore_above": 1024, + "type": "keyword" + }, + "correlationTimespan": { + "ignore_above": 1024, + "type": "keyword" + }, "content": { "type": "text" }, diff --git a/salt/elasticsearch/templates/component/so/detections-alerts-mappings.json b/salt/elasticsearch/templates/component/so/detections-alerts-mappings.json new file mode 100644 index 000000000..5a2abfb7f --- /dev/null +++ b/salt/elasticsearch/templates/component/so/detections-alerts-mappings.json @@ -0,0 +1,149 @@ +{ + "template": { + "mappings": { + "properties": { + "tags": { + "ignore_above": 1024, + "type": "keyword" + }, + "sigma_level": { + "ignore_above": 1024, + "type": "keyword" + }, + "rule": { + "properties": { + "name": { + "ignore_above": 1024, + "type": "keyword" + }, + "uuid": { + "ignore_above": 1024, + "type": "keyword" + }, + "category": { + "ignore_above": 1024, + "type": "keyword" + }, + "product": { + "ignore_above": 1024, + "type": "keyword" + }, + "service": { + "ignore_above": 1024, + "type": "keyword" + }, + "correlation": { + "ignore_above": 1024, + "type": "keyword" + } + } + }, + "event": { + "properties": { + "severity": { + "type": "long" + }, + "severity_label": { + "ignore_above": 1024, + "type": "keyword" + }, + "module": { + "ignore_above": 1024, + "type": "keyword" + }, + "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 + } + } + }, + "event_data": { + "properties": { + "@timestamp": { + "type": "date" + }, + "window_start": { + "type": "date" + }, + "event_count": { + "type": "long" + }, + "value_count": { + "type": "long" + }, + "event_type_count": { + "type": "long" + }, + "value_sum": { + "type": "double" + }, + "value_avg": { + "type": "double" + }, + "value_percentile": { + "type": "double" + }, + "value_median": { + "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" + } + } + } + } + } + }, + "_meta": { + "description": "Fields written by the ElastAlert SecurityOnionESAlerter to logs-detections.alerts-so" + } +} diff --git a/salt/manager/tools/sbin/soup b/salt/manager/tools/sbin/soup index 20c2528dc..5502a378a 100755 --- a/salt/manager/tools/sbin/soup +++ b/salt/manager/tools/sbin/soup @@ -1183,6 +1183,11 @@ up_to_3.4.0() { echo "Removing so-kratos, so-hydra and so-soc so they are recreated on the soauth network." docker rm -f so-kratos so-hydra so-soc >> $SOUP_LOG 2>&1 + # backfill the Sigma rule type on existing detections + mkdir -p /opt/so/conf/soc/migrations + echo "0" > /opt/so/conf/soc/migrations/elastalert-migration-3.4.0 + chown -R socore:socore /opt/so/conf/soc/migrations + for template in so-metrics-logstash.node so-metrics-logstash.stack_monitoring.node; do if ! remove_elasticsearch_index_template "$template" "logstash node and node_cel index patterns reversed"; then FINAL_MESSAGE_QUEUE+=("WARNING: Unable to automatically remove the $template index template. Addon integration templates may fail to load until it is removed:") diff --git a/salt/soc/defaults.yaml b/salt/soc/defaults.yaml index 4eeae5579..f6a97e878 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 @@ -1467,6 +1467,8 @@ soc: - emerging_threats_addon useEsql: false esqlCaseInsensitive: true + esqlQueryDelaySeconds: 30 + esqlCorrelationAllowanceSeconds: 600 elastic: hostUrl: remoteHostUrls: [] @@ -2679,8 +2681,11 @@ soc: query: "so_detection.language:suricata | groupby so_detection.ruleset so_detection.isEnabled | groupby so_detection.category" description: Show all NIDS Detections, which are run with Suricata - name: "Detection Type - Sigma (Elastalert) - All" - query: "so_detection.language:sigma | groupby so_detection.ruleset so_detection.isEnabled | groupby so_detection.category | groupby so_detection.product" + query: "so_detection.language:sigma | groupby so_detection.ruleType | groupby so_detection.ruleset so_detection.isEnabled | groupby so_detection.category | groupby so_detection.product" description: Show all Sigma Detections, which are run with Elastalert + - name: "Detection Type - Sigma (Elastalert) - Correlations" + query: "so_detection.ruleType:correlation | groupby so_detection.correlationType so_detection.isEnabled | groupby so_detection.correlationTimespan | groupby so_detection.ruleset" + description: Show Sigma correlation Detections - name: "Detection Type - YARA (Strelka)" query: "so_detection.language:yara | groupby so_detection.ruleset so_detection.isEnabled" description: Show all YARA detections, which are used by Strelka @@ -2764,7 +2769,7 @@ soc: elastalert: | # This is a Sigma rule template, which uses YAML. Replace all template values with your own values. # The id (UUIDv4) is pregenerated and can safely be used. - # Click "Convert" to convert the Sigma rule to use Security Onion field mappings within an EQL query + # Click "Convert" to convert the Sigma rule to use Security Onion field mappings within a backend query # # Rule Creation Guide: https://github.com/SigmaHQ/sigma/wiki/Rule-Creation-High%E2%80%90Level-Guide # Logsources: https://sigmahq.io/docs/basics/log-sources.html @@ -2794,6 +2799,58 @@ soc: - ' -priv' condition: all of selection_* level: 'high' # info | low | medium | high | critical + elastalert_correlation: | + # Sigma correlation rule; requires ES|QL (useEsql). + # First document: the correlation. Following documents: the rules it references. + # + # Types: event_count, value_count, temporal, value_sum, value_avg, value_median, value_percentile. + # Correlation Guide: https://sigmahq.io/docs/meta/correlations.html + # Logsources: https://sigmahq.io/docs/basics/log-sources.html + + title: 'A Short Capitalized Title With Less Than 50 Characters' + id: [publicId] + status: 'experimental' + description: | + Describe what the correlation finds and, importantly, why relating these + events is more meaningful than either of them alone. + references: + - 'https://local.invalid' + author: '@SecurityOnion' + date: '[today]' + tags: + - detection.threat_hunting + - attack.technique_id + correlation: + type: value_count + rules: + - example_base_rule # the 'name' of the rule below + group-by: + - source.ip + # Xs, Xm, Xh, Xd or Xw. + timespan: 10m + condition: + field: dns.query.name + gte: 40 + falsepositives: + - 'Describe the benign activity that also produces this pattern' + # Placeholders: %count%, %start%, %end%, %duration%, or any field. + summary: '%count% distinct names queried by %source.ip% in %duration%' + level: 'medium' # info | low | medium | high | critical + --- + title: 'Base Event' + # referenced by 'name' (or 'id') + name: example_base_rule + description: 'The single event that the correlation aggregates.' + logsource: + category: network + service: dns + detection: + selection: + dns.query.name|exists: true + condition: selection + # Carried into the alert. + fields: + - dns.query.name assistant: enabled: false investigationPrompt: Investigate Alert ID {socId} diff --git a/salt/soc/files/soc/sigma_pipelines/sigma_esql_pipeline.yml b/salt/soc/files/soc/sigma_pipelines/sigma_esql_pipeline.yml index 18b3a9f63..bc496ef48 100644 --- a/salt/soc/files/soc/sigma_pipelines/sigma_esql_pipeline.yml +++ b/salt/soc/files/soc/sigma_pipelines/sigma_esql_pipeline.yml @@ -475,3 +475,10 @@ transformations: - type: logsource product: kratos +# SOC reads the mapped group-by columns from this output +postprocessing: + - id: esql_correlation_group_by + type: template + template: '{{ {"query": query, "group_by": rule.group_by} | tojson }}' + rule_conditions: + - type: is_sigma_correlation_rule diff --git a/salt/soc/files/soc/sigma_pipelines/sigma_so_pipeline.yml b/salt/soc/files/soc/sigma_pipelines/sigma_so_pipeline.yml index 89b24f53b..38455a5ed 100644 --- a/salt/soc/files/soc/sigma_pipelines/sigma_so_pipeline.yml +++ b/salt/soc/files/soc/sigma_pipelines/sigma_so_pipeline.yml @@ -14,6 +14,15 @@ transformations: - process.args - related.ip - dns.resolved_ip + # always lowercase: matched exactly with the indexed ':' operator + - id: case_sensitive_categorization_fields + type: set_state + key: case_insensitive_exempt_fields + val: + - tags + - event.category + - event.type + - event.kind # Not every source maps .caseless; EQL/ES|QL already match case-insensitively. - id: caseless_to_parent_fields type: field_name_mapping @@ -118,6 +127,9 @@ transformations: valid_hash_algos: ["MD5", "SHA1", "SHA256", "SHA512", "IMPHASH"] field_prefix: "file" drop_algo_prefix: False + # ecs_windows renamed Hashes; pySigma 1.5+ parses only these + field_to_parse: + - winlog.event_data.Hashes field_name_conditions: - type: include_fields fields: diff --git a/salt/soc/merged.map.jinja b/salt/soc/merged.map.jinja index 452fba0b9..4673f1915 100644 --- a/salt/soc/merged.map.jinja +++ b/salt/soc/merged.map.jinja @@ -8,6 +8,7 @@ {% from 'elasticsearch/config.map.jinja' import ELASTICSEARCH_NODES %} {% from 'manager/map.jinja' import MANAGERMERGED %} {% from 'telegraf/map.jinja' import TELEGRAFMERGED %} +{% from 'elastalert/map.jinja' import ELASTALERTMERGED %} {%- set PG_ENTRY = salt['pillar.get']('telegraf:postgres_creds:' ~ grains.id, {}) %} {%- set PG_USER = PG_ENTRY.get('user', '') %} {%- set PG_PASS = PG_ENTRY.get('pass', '') %} @@ -63,6 +64,10 @@ {% do SOCMERGED.config.server.modules.elastalertengine.update({'enabledSigmaRules': SOCMERGED.config.server.modules.elastalertengine.enabledSigmaRules.default}) %} {% endif %} +{# correlation schedules follow ElastAlert's run_every #} +{% set run_every = ELASTALERTMERGED.config.run_every %} +{% do SOCMERGED.config.server.modules.elastalertengine.update({'elastAlertRunEverySeconds': run_every.get('minutes', 0) * 60 + run_every.get('seconds', 0)}) %} + {# set elastalertengine.rulesRepos, strelkaengine.rulesRepos, and suricataengine.rulesetSources based on airgap or not #} {% if GLOBALS.airgap %} {% do SOCMERGED.config.server.modules.elastalertengine.update({'rulesRepos': SOCMERGED.config.server.modules.elastalertengine.rulesRepos.airgap}) %} @@ -80,6 +85,12 @@ {% do SOCMERGED.config.server.update({'airgapEnabled': false}) %} {% endif %} +{# correlation authoring requires ES|QL #} +{% if not SOCMERGED.config.server.modules.elastalertengine.useEsql %} +{% do SOCMERGED.config.server.client.detection.templateDetections.pop('elastalert_correlation', None) %} +{% do SOCMERGED.config.server.client.detections.update({'queries': SOCMERGED.config.server.client.detections.queries | rejectattr('name', 'equalto', 'Detection Type - Sigma (Elastalert) - Correlations') | list}) %} +{% endif %} + {# Define the postgresmetrics module if telegraf is setup to only use Postgres #} {% if TELEGRAFMERGED.output != 'INFLUXDB' and PG_USER and PG_PASS %} {% do SOCMERGED.config.server.modules.update({ diff --git a/salt/soc/soc_soc.yaml b/salt/soc/soc_soc.yaml index df9c142f9..a7f70feeb 100644 --- a/salt/soc/soc_soc.yaml +++ b/salt/soc/soc_soc.yaml @@ -401,15 +401,27 @@ 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." global: True advanced: True forcedType: bool esqlCaseInsensitive: - description: "Match string values case-insensitively when converting Sigma rules. Applies to ES|QL only" + description: "Match string values case-insensitively when converting Sigma rules, and group correlation values regardless of case. Applies to ES|QL only" global: True advanced: True forcedType: bool + esqlQueryDelaySeconds: + description: "Seconds ES|QL rules search behind now, so unsearchable events aren't missed. Delays alerts by the same amount. Set at least the longest index refresh interval. ES|QL only." + global: True + advanced: True + forcedType: int + helpLink: sigma + esqlCorrelationAllowanceSeconds: + description: "Extra seconds of arrivals each correlation run re-reads beyond its timespan, so a burst whose events arrive spread out is still counted together. Correlations below a threshold (lt, lte, eq, neq) and value_avg or value_percentile correlations count only the timespan ending this much plus esqlQueryDelaySeconds ago, so they alert this much later. ES|QL only." + global: True + advanced: True + forcedType: int + helpLink: sigma elastic: index: description: Comma-separated list of indices or index patterns (wildcard "*" supported) that SOC will search for records. 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" | \