This commit is contained in:
defensivedepth committed 2026-10-07 06:28:52 -04:00
1 parent 4244f72d95
commit 506918326d
5 files changed
+147 -108

No files matched your search

@@ -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_<key>"
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_<key>"
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}")
@@ -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)
+1 -1
View File
@@ -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
@@ -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
+2 -2
View File
@@ -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