diff --git a/salt/elastalert/files/modules/so/securityonion-es.py b/salt/elastalert/files/modules/so/securityonion-es.py index 2ee0c0145..f03ed683b 100644 --- a/salt/elastalert/files/modules/so/securityonion-es.py +++ b/salt/elastalert/files/modules/so/securityonion-es.py @@ -6,14 +6,14 @@ # Elastic License 2.0. -from datetime import datetime from time import gmtime, strftime import hashlib import ipaddress import re +import uuid import requests,json from elastalert.alerts import Alerter, DateTimeEncoder -from elastalert.util import EAException, elastalert_logger +from elastalert.util import EAException, elastalert_logger, lookup_es_key, ts_to_dt import urllib3 urllib3.disable_warnings(urllib3.exceptions.InsecureRequestWarning) @@ -26,53 +26,44 @@ class SecurityOnionESAlerter(Alerter): required_options = set(['detection_title', 'sigma_level']) optional_fields = ['sigma_category', 'sigma_product', 'sigma_service', 'sigma_correlation'] - count_labels = { - 'event_count': '%count% events', - 'value_count': '%count% distinct values', - 'event_type_count': '%count% correlated rules matched', - 'value_sum': 'total %count%', - 'value_avg': 'average %count%', - 'value_percentile': 'percentile %count%', - 'value_median': 'median %count%', + # count column and default summary per correlation type + 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'} - @staticmethod - def lookup(doc, dotted): - """ Resolve a dotted path; ES|QL columns arrive nested. """ - node = doc - for part in dotted.split('.'): - if not isinstance(node, dict) or part not in node: - return None - node = node[part] - return node - def query_keys(self): """ compound_query_key holds the list; query_key is flattened to a string. """ return self.rule.get('compound_query_key') or ([self.rule['query_key']] if self.rule.get('query_key') else []) - @staticmethod - def is_correlation(match): - """ Only correlation rows carry window_start, from the ES|QL stats. """ - return 'window_start' in match + def is_correlation(self): + return bool(self.rule.get('sigma_correlation')) def alert_id(self, match): - """ Stable id: window end + group values for correlations, source _id otherwise. """ - if self.is_correlation(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"|{self.lookup(match, k)}" for k in self.query_keys()) - key = f"{self.rule['detection_public_id']}|{self.to_dt(match['@timestamp']).isoformat()}{values}" + 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: - key = f"{self.rule['detection_public_id']}|{match.get('_id')}" + return uuid.uuid4().hex 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()) + 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'): @@ -98,7 +89,10 @@ class SecurityOnionESAlerter(Alerter): bucket = self.related_bucket(key) if not bucket: continue - value = self.lookup(match, key) + # a lowercased group carries its original spellings, which keyword searches need + 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 @@ -133,21 +127,16 @@ class SecurityOnionESAlerter(Alerter): return f"{seconds // size} {unit}s" return f"{seconds} second{'' if seconds == 1 else 's'}" - @staticmethod - def to_dt(value): - # ES|QL gives ISO strings; ElastAlert parses @timestamp, except on a retried alert. - return value if isinstance(value, datetime) else datetime.fromisoformat(value) - def summary(self, match): """ One-line correlation summary; None for single-event rules. """ - if not self.is_correlation(match): + if not self.is_correlation(): return None - start = self.to_dt(match['window_start']) - end = self.to_dt(match['@timestamp']) - name = next((f for f in self.count_labels if f in match), None) + 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(name)), + '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())), @@ -156,12 +145,12 @@ class SecurityOnionESAlerter(Alerter): template = self.rule.get('summary_template') if not template: groups = ', '.join(f"{k} %{k}%" for k in self.query_keys()) - template = f"{self.count_labels.get(name, '%count%')}{' for ' + groups if groups else ''} in %duration%" + 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 = self.lookup(match, 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) @@ -169,76 +158,81 @@ class SecurityOnionESAlerter(Alerter): def alert(self, matches): for match in matches: - timestamp = strftime("%Y-%m-%d"'T'"%H:%M:%S"'.000Z', gmtime()) - # Start building the rule dict - rule_info = { - "name": self.rule['detection_title'], - "uuid": self.rule['detection_public_id'] + self.write(self.alert_id(match), self.payload(match)) + + def payload(self, match): + timestamp = strftime("%Y-%m-%d"'T'"%H:%M:%S"'.000Z', gmtime()) + # 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] + + 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"], + "rule": rule_info, + "event": event_info, + "sigma_level": self.rule['sigma_level'], + "event_data": self.event_data(match), + "@timestamp": timestamp + } + + keys = self.query_keys() + if keys and self.is_correlation(): + payload["labels"] = { + "correlation_group_by": ', '.join(keys), + "correlation_group": self.group(match), } + related = self.related(match) + if related: + payload["related"] = related - # 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] + return payload - 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"], - "rule": rule_info, - "event": event_info, - "sigma_level": self.rule['sigma_level'], - "event_data": self.event_data(match), - "@timestamp": timestamp - } - - keys = self.query_keys() - if keys and self.is_correlation(match): - 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}") - response = self.put_alert(url, payload) + def write(self, alert_id, payload): + # _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}") + response = self.put_alert(url, payload) + if response.status_code == 400: + # mapping rejections come from event_data; retry with it as unindexed text + rejection = response.text[:500] + response = self.put_alert(url, self.without_event_data(payload, rejection)) if response.status_code == 400: - # 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]}") + 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) + 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) + 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) + return requests.put(url, data=json.dumps(payload, cls=DateTimeEncoder), headers={"Content-Type": "application/json"}, + verify=False, auth=creds, timeout=self.rule.get('es_conn_timeout', 20)) except requests.RequestException as e: raise EAException(f"Unable to write alert: {e}") diff --git a/salt/elastalert/files/modules/so/securityonion-es_test.py b/salt/elastalert/files/modules/so/securityonion-es_test.py index 5cc2c5043..0e56c45ab 100644 --- a/salt/elastalert/files/modules/so/securityonion-es_test.py +++ b/salt/elastalert/files/modules/so/securityonion-es_test.py @@ -30,12 +30,24 @@ except ImportError: 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) + 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 package = types.ModuleType('elastalert') package.alerts = alerts package.util = util @@ -56,9 +68,12 @@ BASE_RULE = { '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 { @@ -82,6 +97,7 @@ class TestSecurityOnionESAlerter(unittest.TestCase): with patch.object(es.requests, 'put', return_value=response) as put: alerter.alert([match]) self.assertEqual(put.call_count, 1) + self.assertEqual(put.call_args.kwargs['timeout'], 55) return json.loads(put.call_args.kwargs['data']), put.call_args.args[0] def test_compound_query_key_left_out_of_event_data(self): @@ -131,9 +147,21 @@ class TestSecurityOnionESAlerter(unittest.TestCase): 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, as from an override, does not make a correlation - rule = dict(BASE_RULE, query_key='user.name') + 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': 'josh'}} @@ -148,6 +176,8 @@ class TestSecurityOnionESAlerter(unittest.TestCase): self.assertEqual(payload['event_data'], dict(match, **{'@timestamp': '2026-09-30T16:50:00+00:00'})) self.assertTrue(url.endswith('/' + 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) @@ -164,6 +194,14 @@ class TestSecurityOnionESAlerter(unittest.TestCase): 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') @@ -202,7 +240,7 @@ class TestSecurityOnionESAlerter(unittest.TestCase): 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]) + payloads = self.send_responses(PLAIN_RULE, match, [rejected, rejected]) self.assertEqual(len(payloads), 2) diff --git a/salt/soc/defaults.yaml b/salt/soc/defaults.yaml index 005c2e9eb..2f8af8713 100644 --- a/salt/soc/defaults.yaml +++ b/salt/soc/defaults.yaml @@ -2803,7 +2803,7 @@ soc: # Sigma correlation rule; requires ES|QL (useEsql). # The first document is the correlation (id, title, severity); the documents after it are the rules it refers to, all in this Detection. # - # Supported types: event_count, value_count, temporal, value_sum, value_avg, value_percentile, value_median. + # Supported types: event_count, value_count, temporal, value_sum, value_avg, value_percentile. # Correlation Guide: https://sigmahq.io/docs/meta/correlations.html # Logsources: https://sigmahq.io/docs/basics/log-sources.html 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..ace10725f 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 a correlation's group-by columns, as mapped by the pipelines, 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/soc_soc.yaml b/salt/soc/soc_soc.yaml index fe2cad90a..2e4438d96 100644 --- a/salt/soc/soc_soc.yaml +++ b/salt/soc/soc_soc.yaml @@ -406,7 +406,7 @@ soc: 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 @@ -417,7 +417,7 @@ soc: 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. ES|QL only." + 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