Compare commits

...
Author SHA1 Message Date
defensivedepth 98ffb6fa00 Initial Correlations support 2026-09-29 17:05:20 -04:00
11 changed files with 328 additions and 8 deletions

No files matched your search

+1 -1
View File
@@ -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
@@ -6,9 +6,13 @@
# Elastic License 2.0.
from datetime import datetime
from time import gmtime, strftime
import hashlib
import re
import requests,json
from elastalert.alerts import Alerter
from elastalert.alerts import Alerter, DateTimeEncoder
from elastalert.util import EAException, elastalert_logger
import urllib3
urllib3.disable_warnings(urllib3.exceptions.InsecureRequestWarning)
@@ -19,7 +23,99 @@ class SecurityOnionESAlerter(Alerter):
"""
required_options = set(['detection_title', 'sigma_level'])
optional_fields = ['sigma_category', 'sigma_product', 'sigma_service']
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%',
}
placeholder = re.compile(r'%([^%\s]+)%')
@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 [])
def alert_id(self, match):
""" Stable id: window + group values for correlations, source _id otherwise. """
keys = self.query_keys()
if keys:
values = '|'.join(str(self.lookup(match, k)) for k in keys)
key = f"{self.rule['detection_public_id']}|{self.to_dt(match['@timestamp']).isoformat()}|{values}"
else:
key = f"{self.rule['detection_public_id']}|{match.get('_id')}"
return hashlib.sha256(key.encode('utf-8')).hexdigest()
@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'}"
@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 'window_start' not in match:
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)
values = {
'count': self.format_count(match.get(name)),
'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"{self.count_labels.get(name, '%count%')}{' for ' + groups if groups else ''} in %duration%"
def fill(m):
if m[1] in values:
return values[m[1]]
value = self.lookup(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:
@@ -42,6 +138,10 @@ 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
# Construct the payload with the conditional rule_info
payload = {
"tags": "alert",
@@ -56,8 +156,21 @@ class SecurityOnionESAlerter(Alerter):
"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)
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}")
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
if response.status_code != 409 and not response.ok:
raise EAException(f"Unable to write alert: {response.status_code} {response.text[:500]}")
def get_info(self):
return {'type': 'SecurityOnionESAlerter'}
+1 -1
View File
@@ -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:
+1
View File
@@ -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
@@ -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"
},
@@ -0,0 +1,106 @@
{
"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"
},
"summary": {
"type": "match_only_text",
"fields": {
"keyword": {
"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"
}
}
},
"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"
}
}
}
}
}
},
"_meta": {
"description": "Fields written by the ElastAlert SecurityOnionESAlerter to logs-detections.alerts-so"
}
}
+5
View File
@@ -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
# Extract the Sigma rule type for existing detections (Single vs. Correlation)
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
INSTALLEDVERSION=3.4.0
}
+59 -2
View File
@@ -1467,6 +1467,8 @@ soc:
- emerging_threats_addon
useEsql: false
esqlCaseInsensitive: true
esqlQueryDelaySeconds: 30
esqlCorrelationAllowanceSeconds: 600
elastic:
hostUrl:
remoteHostUrls: []
@@ -2671,8 +2673,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
@@ -2756,7 +2761,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
@@ -2786,6 +2791,58 @@ soc:
- ' -priv'
condition: all of selection_*
level: 'high' # info | low | medium | high | critical
elastalert_correlation: |
# 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.
# 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 # matches the 'name' of the document 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%, and group-by or field values.
summary: '%count% distinct names queried by %source.ip% in %duration%'
level: 'medium' # info | low | medium | high | critical
---
title: 'Base Event'
# The correlation refers to this rule by 'name' (or by '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}
@@ -14,6 +14,15 @@ transformations:
- process.args
- related.ip
- dns.resolved_ip
# Always lowercase, so matched exactly; the backend can then use the indexed ':' operator.
- id: case_sensitive_categorization_fields
type: set_state
key: case_insensitive_exempt_fields
val:
- tags
- event.category
- event.type
- event.kind
- id: esql_default_index
type: set_state
key: index
+5
View File
@@ -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}) %}
+12
View File
@@ -409,6 +409,18 @@ soc:
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. 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.