Additional ESQL tweaks

This commit is contained in:
defensivedepth committed 2026-10-01 11:32:42 -04:00
1 parent 98ffb6fa00
commit 3d4f53b741
7 files changed
+407 -39

No files matched your search

@@ -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'}
@@ -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()
@@ -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
}
}
]
}
@@ -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"
}
}
}
}
}
+3 -3
View File
@@ -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
+1 -1
View File
@@ -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
+1 -1
View File
@@ -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" | \